Skip to content

Commit 735f70d

Browse files
committed
improvement(knowledge): scan each reconciliation window as one ordered, limited walk of the v2 index
1 parent e875419 commit 735f70d

6 files changed

Lines changed: 129 additions & 125 deletions

File tree

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

Lines changed: 28 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ import {
1919
user,
2020
workspace,
2121
} from '@sim/db/schema'
22-
import { toNumberOrNull, toStringOrNull } from '@sim/utils/coerce'
22+
import { toNumberOrNull } 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'
@@ -963,9 +963,7 @@ describe('durable source and member cycles in PostgreSQL', () => {
963963
await db.insert(document).values(rows.slice(offset, offset + 1_000))
964964
await db.execute(sql`ANALYZE document`)
965965
}
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)
966+
const isWindowScan = (query: string) => /^\s*with "document" as materialized/i.test(query)
969967
/** Runs the pass, returning the window statements it issued outside transactions. */
970968
const reconcile = async () => {
971969
const hardDelete = vi.spyOn(documentService, 'hardDeleteDocuments')
@@ -1001,19 +999,20 @@ describe('durable source and member cycles in PostgreSQL', () => {
1001999
hardDeleted: hardDelete.mock.calls.flatMap(([batch]) => batch),
10021000
walked: statements.mock.calls
10031001
.map(([query, params]) => ({ query, params }))
1004-
.filter(({ query }) => isWindowBound(query) || isWindowPage(query)),
1002+
.filter(({ query }) => isWindowScan(query)),
10051003
}
10061004
} finally {
10071005
statements.mockRestore()
10081006
hardDelete.mockRestore()
10091007
}
10101008
}
10111009
/**
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.
1010+
* Plans a window statement as a large table would, on an index rather than a sequential
1011+
* scan, from freshly analyzed statistics, returning the document rows it read.
10141012
*/
10151013
const explainWalk = async (query: string, params: Parameters<typeof db.$client.unsafe>[1]) => {
10161014
const [explained] = await db.$client.begin(async (tx) => {
1015+
await tx.unsafe('ANALYZE document')
10171016
await tx.unsafe('SET LOCAL enable_seqscan = off')
10181017
return tx.unsafe(`EXPLAIN (ANALYZE, FORMAT JSON) ${query}`, params)
10191018
})
@@ -1022,28 +1021,30 @@ describe('durable source and member cycles in PostgreSQL', () => {
10221021
return [plan, ...toArray(plan.Plans).flatMap(nodes)]
10231022
}
10241023
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) =>
1024+
return plan
1025+
.filter((node) => node['Relation Name'] === 'document')
1026+
.reduce(
1027+
(total, node) =>
10291028
total +
1030-
((toNumberOrNull(plan['Actual Rows']) ?? 0) +
1031-
(toNumberOrNull(plan['Rows Removed by Filter']) ?? 0)) *
1032-
(toNumberOrNull(plan['Actual Loops']) ?? 1),
1029+
((toNumberOrNull(node['Actual Rows']) ?? 0) +
1030+
(toNumberOrNull(node['Rows Removed by Filter']) ?? 0)) *
1031+
(toNumberOrNull(node['Actual Loops']) ?? 1),
10331032
0
1034-
),
1035-
indexes: plan.map((node) => toStringOrNull(node['Index Name'])).filter(Boolean),
1036-
}
1033+
)
10371034
}
1038-
/** Every window statement reads through the connector's v2 id range, at most one window of it. */
1035+
/**
1036+
* Every window statement reads at most one window of documents, whichever connector index
1037+
* the planner walks: a read of the whole connector, or of every tombstone it has, is what
1038+
* must never happen.
1039+
*/
10391040
const expectBounded = async (
10401041
walked: { query: string; params: Parameters<typeof db.$client.unsafe>[1] }[]
10411042
) => {
10421043
expect(walked.length).toBeGreaterThan(0)
10431044
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)
1045+
expect(await explainWalk(query, params), query).toBeLessThanOrEqual(
1046+
RECONCILIATION_WINDOW_SIZE
1047+
)
10471048
}
10481049
}
10491050

@@ -1079,7 +1080,7 @@ describe('durable source and member cycles in PostgreSQL', () => {
10791080
await expectBounded(walked)
10801081
}, 120_000)
10811082

1082-
it('bounds a dense window once and pages it on the id range, not the tombstone index', async () => {
1083+
it('scans a dense window once and never reads every tombstone of the connector', async () => {
10831084
/**
10841085
* Three hard-delete pages of absent tombstones open the first window; the rest of the
10851086
* connector is tombstones the listing still sees, which the hard walk must pass over.
@@ -1105,11 +1106,13 @@ describe('durable source and member cycles in PostgreSQL', () => {
11051106
.from(document)
11061107
.where(eq(document.connectorId, connectorId))
11071108
).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(
1109+
/**
1110+
* The only walk is the hard one: one statement per window, the dense first one included,
1111+
* then the tail, however many pages of matches a window holds.
1112+
*/
1113+
expect(walked).toHaveLength(
11101114
Math.floor((dense.length + seenTombstones.length) / RECONCILIATION_WINDOW_SIZE) + 1
11111115
)
1112-
expect(walked.filter(({ query }) => isWindowPage(query)).length).toBeGreaterThan(3)
11131116
/** `deleted_at < $1` implies the tombstone index, which would read every connector tombstone. */
11141117
await expectBounded(walked)
11151118
}, 120_000)

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

Lines changed: 7 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ import { document, embedding, knowledgeConnector, user, workspace } from '@sim/d
44
import { createLogger } from '@sim/logger'
55
import { toStringOrNull } from '@sim/utils/coerce'
66
import { toArray, toRecord } from '@sim/utils/object'
7-
import { and, asc, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
7+
import { and, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
88
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
99
import { z } from 'zod'
1010
import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution'
@@ -325,25 +325,12 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
325325
AND (source_seen_at IS NULL OR source_seen_at < ${startedAt.toISOString()}::timestamp) AND cardinality(acl) > 0 LIMIT 500`
326326
)
327327
const owned = sql`connector_id = ${ids.connectorId} AND user_excluded = false AND archived_at IS NULL`
328-
const windowEnd = db
329-
.select({ id: document.id })
330-
.from(document)
331-
.where(owned)
332-
.orderBy(asc(document.id))
333-
.limit(1)
334-
.offset(RECONCILIATION_WINDOW_SIZE - 1)
335-
await explain('reconciliation.window.end', windowEnd.getSQL())
336-
expect(planIndexNames('reconciliation.window.end')).toContain(
337-
'doc_connector_reconciliation_v2_idx'
338-
)
339-
const [end] = await windowEnd
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)
346-
expect(planIndexNames('reconciliation.window.page')).toContain(
328+
const windowScan = sql`WITH document AS MATERIALIZED (
329+
SELECT id, source_seen_at, acl FROM document WHERE ${owned} ORDER BY id LIMIT ${RECONCILIATION_WINDOW_SIZE})
330+
SELECT (SELECT max(id) FROM document) AS last, json_agg(id ORDER BY id) AS ids FROM document
331+
WHERE (source_seen_at IS NULL OR source_seen_at < ${startedAt.toISOString()}::timestamp) AND cardinality(acl) > 0`
332+
await explain('reconciliation.window.scan', windowScan)
333+
expect(planIndexNames('reconciliation.window.scan')).toContain(
347334
'doc_connector_reconciliation_v2_idx'
348335
)
349336
const outcome = await measure('reconciliation.actual', async () =>

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -339,8 +339,8 @@ 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; its page, read through `db.execute`, is empty. */
343-
queueTableRows(schemaMock.document, [])
342+
/** The resurrection walk's only window, read through `db.execute`, is empty. */
343+
dbChainMockFns.execute.mockResolvedValueOnce([{ size: 0, last: null, ids: [], tombstoned: [] }])
344344
queueTableRows(schemaMock.document, [{ id: 'd-1' }])
345345
vi.mocked(hardDeleteDocuments).mockRejectedValueOnce(
346346
new ConnectorSyncDeletionGuardError('lease reclaimed')
Lines changed: 60 additions & 60 deletions
Original file line numberDiff line numberDiff line change
@@ -1,18 +1,17 @@
11
import { db } from '@sim/db'
22
import { document } from '@sim/db/schema'
3-
import { and, asc, eq, gt, isNull, lte, type SQL, sql } from 'drizzle-orm'
3+
import { and, asc, eq, gt, isNull, type SQL, sql } from 'drizzle-orm'
44

55
/**
66
* Ids of `doc_connector_reconciliation_v2_idx` one reconciliation statement may scan. The rows a
77
* walk looks for can be rare and late in id order, so a page bounded only by its matches could
8-
* read, and heap-fetch, a whole connector in one statement. Each page instead scans at most this
9-
* many ids: a few hundred milliseconds cold, and few enough round trips that a large connector
10-
* is walked in hundreds of pages rather than thousands.
8+
* read, and heap-fetch, a whole connector in one statement. Each window instead scans at most
9+
* this many ids: a few hundred milliseconds cold, and few enough round trips that a large
10+
* connector is walked in hundreds of statements rather than thousands.
1111
*/
1212
export const RECONCILIATION_WINDOW_SIZE = 5_000
1313

14-
/** A type alias, not an interface, so it satisfies `db.execute`'s row-record constraint. */
15-
type ReconciliationRow = {
14+
interface ReconciliationRow {
1615
id: string
1716
tombstoned: boolean
1817
}
@@ -31,40 +30,27 @@ export interface ReconciliationWalk {
3130
onPage: (rows: ReconciliationRow[]) => Promise<void>
3231
}
3332

34-
function ownedAfter(connectorId: string, afterId: string | undefined) {
35-
return and(
36-
eq(document.connectorId, connectorId),
37-
eq(document.userExcluded, false),
38-
isNull(document.archivedAt),
39-
afterId ? gt(document.id, afterId) : undefined
40-
)
41-
}
42-
43-
/** The last id of the next full window after `afterId`, or undefined when less than a window remains. */
44-
async function windowEnd(connectorId: string, afterId: string | undefined) {
45-
const [row] = await db
46-
.select({ id: document.id })
47-
.from(document)
48-
.where(ownedAfter(connectorId, afterId))
49-
.orderBy(asc(document.id))
50-
.limit(1)
51-
.offset(RECONCILIATION_WINDOW_SIZE - 1)
52-
return row?.id
33+
/** A type alias, not an interface, so it satisfies `db.execute`'s row-record constraint. */
34+
type WindowScan = {
35+
size: number
36+
last: string | null
37+
ids: string[]
38+
tombstoned: boolean[]
5339
}
5440

5541
/**
56-
* The first `pageSize` rows of the window `(afterId, endId]` matching `condition`. The window is
57-
* read on its own, as a materialized CTE that shadows `document`, so the planner reaches it only
58-
* through the v2 id range: a condition that implies a narrower partial index, such as the
59-
* tombstone index for `deleted_at < $1`, would otherwise draw a bitmap over every tombstone the
60-
* connector has.
42+
* Reads the next window after `afterId` and the ids in it matching `condition`, in one statement.
43+
* The matches come back as JSON, which the driver decodes without the array type lookup.
44+
* The window is the first {@link RECONCILIATION_WINDOW_SIZE} owned documents in id order, a
45+
* materialized CTE that shadows `document`: the `ORDER BY id LIMIT` makes it an ordered walk of
46+
* the v2 index that stops at the window's size. A range bounded only by `id` is not enough, as
47+
* the planner cannot tell that it covers one window of the connector and may instead read every
48+
* document of the connector through another connector index; and a condition that implies a
49+
* narrower partial index, such as the tombstone index for `deleted_at < $1`, could draw a bitmap
50+
* over every tombstone the connector has.
6151
*/
62-
async function windowPage(
63-
walk: ReconciliationWalk,
64-
afterId: string | undefined,
65-
endId: string | undefined
66-
): Promise<ReconciliationRow[]> {
67-
const rows = db
52+
async function scanWindow(walk: ReconciliationWalk, afterId: string | undefined) {
53+
const window = db
6854
.select({
6955
id: document.id,
7056
connectorId: document.connectorId,
@@ -76,42 +62,56 @@ async function windowPage(
7662
contentHash: document.contentHash,
7763
})
7864
.from(document)
79-
.where(and(ownedAfter(walk.connectorId, afterId), endId ? lte(document.id, endId) : undefined))
80-
return db.execute<ReconciliationRow>(sql`
81-
WITH ${document} AS MATERIALIZED (${rows})
82-
SELECT ${document.id} AS "id", ${document.deletedAt} IS NOT NULL AS "tombstoned"
83-
FROM ${document}
84-
WHERE ${walk.condition ?? sql`true`}
85-
ORDER BY ${document.id}
86-
LIMIT ${walk.pageSize}`)
65+
.where(
66+
and(
67+
eq(document.connectorId, walk.connectorId),
68+
eq(document.userExcluded, false),
69+
isNull(document.archivedAt),
70+
afterId ? gt(document.id, afterId) : undefined
71+
)
72+
)
73+
.orderBy(asc(document.id))
74+
.limit(RECONCILIATION_WINDOW_SIZE)
75+
const [scan] = await db.execute<WindowScan>(sql`
76+
WITH ${document} AS MATERIALIZED (${window}),
77+
matched AS (
78+
SELECT ${document.id} AS id, ${document.deletedAt} IS NOT NULL AS tombstoned
79+
FROM ${document}
80+
WHERE ${walk.condition ?? sql`true`}
81+
)
82+
SELECT
83+
(SELECT count(*)::int FROM ${document}) AS "size",
84+
(SELECT max(${document.id}) FROM ${document}) AS "last",
85+
coalesce(json_agg(id ORDER BY id), '[]') AS "ids",
86+
coalesce(json_agg(tombstoned ORDER BY id), '[]') AS "tombstoned"
87+
FROM matched`)
88+
return scan
8789
}
8890

8991
/**
90-
* Walks a connector's owned documents in id order, one bounded window per statement, handing
91-
* each page of the rows matching `condition` to `onPage`. A window is bounded once and paged
92-
* until its matches run out, so a dense window costs one bound however many pages it holds. A
93-
* window without matches still advances the walk, which ends only when the connector's tail (a
94-
* window with no end) is exhausted. Returns false when the deadline stopped it.
92+
* Walks a connector's owned documents in id order, one window per statement, handing the rows
93+
* matching `condition` to `onPage` in pages of at most `pageSize`. A window without matches still
94+
* advances the walk, which ends after the first window shorter than a full one. Returns false when
95+
* the deadline stopped it.
9596
*/
9697
export async function walkReconciliationWindows(walk: ReconciliationWalk): Promise<boolean> {
9798
let after: string | undefined
98-
let window: { end: string | undefined } | undefined
9999
for (;;) {
100100
if (Date.now() >= walk.deadlineAt) return false
101101
await walk.beforePage()
102102
if (Date.now() >= walk.deadlineAt) return false
103-
window ??= { end: await windowEnd(walk.connectorId, after) }
104-
const { end } = window
105-
const rows = await windowPage(walk, after, end)
106-
if (rows.length > 0) {
103+
const { size, last, ids, tombstoned } = await scanWindow(walk, after)
104+
for (let offset = 0; offset < ids.length; offset += walk.pageSize) {
105+
if (offset > 0) await walk.beforePage()
107106
/** Materialized ids are acted on only inside the budget, so a late page is left for the next run. */
108107
if (Date.now() >= walk.deadlineAt) return false
109-
await walk.onPage(rows)
108+
await walk.onPage(
109+
ids
110+
.slice(offset, offset + walk.pageSize)
111+
.map((id, index) => ({ id, tombstoned: tombstoned[offset + index] }))
112+
)
110113
}
111-
if (rows.length === walk.pageSize) after = rows[rows.length - 1].id
112-
else if (end) {
113-
after = end
114-
window = undefined
115-
} else return true
114+
if (size < RECONCILIATION_WINDOW_SIZE || !last) return true
115+
after = last
116116
}
117117
}

‎apps/sim/lib/knowledge/connectors/sync-content-pass.test.ts‎

Lines changed: 17 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -47,9 +47,21 @@ uploadsMetadataMockFns.mockInsertImmutableFileMetadata.mockImplementation(
4747
}
4848
)
4949

50-
/** Reconciliation window pages read through `db.execute`; queued here in call order. */
50+
/** Reconciliation window scans read through `db.execute`; queued here in call order. */
5151
const windowPages: unknown[][] = []
5252

53+
/** The one-row result of a reconciliation window scan shorter than a full window. */
54+
function windowScan(rows: { id: string; tombstoned: boolean }[]) {
55+
return [
56+
{
57+
size: rows.length,
58+
last: rows.at(-1)?.id ?? null,
59+
ids: rows.map((row) => row.id),
60+
tombstoned: rows.map((row) => row.tombstoned),
61+
},
62+
]
63+
}
64+
5365
function isWindowPage(query: unknown) {
5466
const { toSQL } = toRecord(query)
5567
return typeof toSQL === 'function' && String(toSQL().sql).includes('MATERIALIZED')
@@ -233,13 +245,9 @@ describe('completed listing removal counts', () => {
233245
hardCount: hard.length,
234246
},
235247
])
236-
/**
237-
* A walk whose count is zero never reads a page. Each walk here fits in one partial window,
238-
* so it reads the absent window bound, then its only page.
239-
*/
248+
/** A walk whose count is zero never scans. Each walk here fits in one partial window. */
240249
if (options.revoked?.length) {
241-
queueTableRows(schemaMock.document, [])
242-
windowPages.push(options.revoked)
250+
windowPages.push(windowScan(options.revoked))
243251
/** Each window reads what still grants someone; pages are bounded by their chunks' rows. */
244252
const granting = options.revoked.map(({ id }) => ({ id, chunkCount: 10 }))
245253
queueTableRows(schemaMock.document, granting)
@@ -254,16 +262,14 @@ describe('completed listing removal counts', () => {
254262
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'connector' }])
255263
}
256264
if (!options.fullSync && soft.length) {
257-
queueTableRows(schemaMock.document, [])
258-
windowPages.push(soft)
265+
windowPages.push(windowScan(soft))
259266
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'connector' }])
260267
dbChainMockFns.returning.mockResolvedValueOnce(
261268
options.updated ?? soft.map(({ id }) => ({ id }))
262269
)
263270
}
264271
if (hard.length) {
265-
queueTableRows(schemaMock.document, [])
266-
windowPages.push(hard)
272+
windowPages.push(windowScan(hard))
267273
}
268274
const result: SyncResult = {
269275
docsAdded: 0,

0 commit comments

Comments
 (0)