Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -1156,6 +1156,13 @@ describe('connector lease ACL pages in PostgreSQL', () => {
.update(document)
.set({ sourceSeenAt: sql`now() + interval '1 day'` })
.where(eq(document.connectorId, members.connectorId))
const seenStamps = () =>
db
.select({ id: document.id, sourceSeenAt: document.sourceSeenAt })
.from(document)
.where(eq(document.connectorId, members.connectorId))
.orderBy(document.id)
const seenBefore = await seenStamps()
provider.list.mockResolvedValue({
documents: seeded.map((row) => ({
externalId: row.externalId,
Expand Down Expand Up @@ -1186,6 +1193,8 @@ describe('connector lease ACL pages in PostgreSQL', () => {
expect(perTransaction.reduce((total, writes) => total + writes, 0)).toBe(2 * DOCUMENTS)
expect(Math.max(...perTransaction)).toBeLessThanOrEqual(PAGE)
await expectBounded(members.connectorId)
/** A row already stamped at or after this run's start is not rewritten by the seen stamp. */
expect(await seenStamps()).toEqual(seenBefore)
})

it('rematerialises what a change feed withdrew one page per lease transaction', async () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,9 +19,11 @@ import {
user,
workspace,
} from '@sim/db/schema'
import { toNumberOrNull } from '@sim/utils/coerce'
import { generateId } from '@sim/utils/id'
import { toArray, toRecord } from '@sim/utils/object'
import { and, eq, inArray, sql } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'

const fixture = vi.hoisted(() => ({
storageRoot: '',
Expand Down Expand Up @@ -92,6 +94,7 @@ import {
executeMemberSync,
resumeMembershipRewrites,
} from '@/lib/knowledge/connectors/member-sync-engine'
import { RECONCILIATION_WINDOW_SIZE } from '@/lib/knowledge/connectors/reconciliation-window'
import { runConnectorContentPass } from '@/lib/knowledge/connectors/sync-content-pass'
import { executeSync } from '@/lib/knowledge/connectors/sync-engine'
import {
Expand Down Expand Up @@ -898,6 +901,220 @@ describe('durable source and member cycles in PostgreSQL', () => {
}
})

describe('reconciliation windows', () => {
const UNRELATED_FILENAME = 'reconciliation-window-unrelated'
const acl = [`u:${ids.aliceId}@fixture.test`]
let connectorId = ''
let runId = ''
let startedAt = new Date()
beforeEach(() => {
connectorId = generateId()
runId = generateId()
startedAt = new Date()
})
afterEach(async () => {
await db.delete(document).where(eq(document.filename, UNRELATED_FILENAME))
await db.delete(document).where(eq(document.connectorId, connectorId))
await db.delete(knowledgeConnector).where(eq(knowledgeConnector.id, connectorId))
})
const row = (id: string) => ({
id,
knowledgeBaseId: ids.knowledgeBaseId,
connectorId,
externalId: id,
filename: id,
fileUrl: '',
fileSize: 0,
mimeType: 'text/plain',
processingStatus: 'completed',
contentHash: `hash-${id}`,
acl,
aclVerifiedAt: startedAt,
})
/**
* Seeds a connector whose listing already completed, so the pass goes straight to
* reconciliation. Other documents spread across the id space keep the connector the minority
* it is at scale; a connector that is nearly the whole table may be walked through the
* primary key instead, one row at a time all the same, reading the other rows between its
* own.
*/
const seed = async (rows: ReturnType<typeof row>[], listedCount: number) => {
await db.execute(sql`
INSERT INTO ${document} (id, knowledge_base_id, filename, file_url, file_size, mime_type, processing_status)
SELECT md5(${connectorId} || g), ${ids.knowledgeBaseId}, ${UNRELATED_FILENAME}, '', 0, 'text/plain', 'completed'
FROM generate_series(1, ${rows.length}) g`)
const checkpoint = beginListingCheckpoint({
fingerprint: listingFingerprint({ connectorId }),
generationId: runId,
startedAt,
})
checkpoint.complete = true
checkpoint.listedCount = listedCount
await db.insert(knowledgeConnector).values({
id: connectorId,
knowledgeBaseId: ids.knowledgeBaseId,
connectorType: 'google_drive',
sourceConfig: {},
accessMode: 'admin',
status: 'syncing',
syncLockToken: runId,
listingCheckpoint: checkpoint,
})
for (let offset = 0; offset < rows.length; offset += 1_000)
await db.insert(document).values(rows.slice(offset, offset + 1_000))
await db.execute(sql`ANALYZE document`)
}
const isWindowScan = (query: string) => /^\s*with "document" as materialized/i.test(query)
/** Runs the pass, returning the window statements it issued outside transactions. */
const reconcile = async () => {
const hardDelete = vi.spyOn(documentService, 'hardDeleteDocuments')
const statements = vi.spyOn(db.$client, 'unsafe')
try {
const [connector] = await db
.select()
.from(knowledgeConnector)
.where(eq(knowledgeConnector.id, connectorId))
const stats = result()
const pass = await runConnectorContentPass({
connectorId,
connector,
connectorConfig: CONNECTOR_REGISTRY.google_drive,
sourceConfig: {},
syncContext: {},
kbOwner: { userId: ids.aliceId, workspaceId: ids.workspaceId },
billingAttribution: billing,
result: stats,
lease: createContentSyncLease(connectorId, runId),
leaseKind: 'content',
runId,
fingerprint: listingFingerprint({ connectorId }),
documentAccess: 'admin',
getAccessToken: async () => 'fixture',
hydration: { getDocument: fixture.get },
forceRehydrate: false,
deadlineAt: Date.now() + 60_000,
})
return {
pass,
stats,
hardDeleted: hardDelete.mock.calls.flatMap(([batch]) => batch),
walked: statements.mock.calls
.map(([query, params]) => ({ query, params }))
.filter(({ query }) => isWindowScan(query)),
}
} finally {
statements.mockRestore()
hardDelete.mockRestore()
}
}
/** Plans a window statement from freshly analyzed statistics, returning the document rows it read. */
const explainWalk = async (query: string, params: Parameters<typeof db.$client.unsafe>[1]) => {
const [explained] = await db.$client.begin(async (tx) => {
await tx.unsafe('ANALYZE document')
return tx.unsafe(`EXPLAIN (ANALYZE, FORMAT JSON) ${query}`, params)
})
const nodes = (node: unknown): Record<string, unknown>[] => {
const plan = toRecord(node)
return [plan, ...toArray(plan.Plans).flatMap(nodes)]
}
const plan = nodes(toRecord(toArray(explained['QUERY PLAN'])[0]).Plan)
return plan
.filter((node) => node['Relation Name'] === 'document')
.reduce(
(total, node) =>
total +
((toNumberOrNull(node['Actual Rows']) ?? 0) +
(toNumberOrNull(node['Rows Removed by Filter']) ?? 0)) *
(toNumberOrNull(node['Actual Loops']) ?? 1),
0
)
}
/**
* Every window statement reads at most one window of documents, whichever connector index
* the planner walks: a read of the whole connector, or of every tombstone it has, is what
* must never happen.
*/
const expectBounded = async (
walked: { query: string; params: Parameters<typeof db.$client.unsafe>[1] }[]
) => {
expect(walked.length).toBeGreaterThan(0)
for (const { query, params } of walked) {
expect(await explainWalk(query, params), query).toBeLessThanOrEqual(
RECONCILIATION_WINDOW_SIZE
)
}
}

it('bounds every page by the ids it scans when absence is rare and late in id order', async () => {
/** Ids sort present rows first, so each absent row is found only after three full windows. */
const present = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
...row(`${connectorId}-a-${String(index).padStart(6, '0')}`),
sourceSeenAt: startedAt,
}))
const absent = Array.from({ length: 3 }, (_, index) => ({
...row(`${connectorId}-z-live-${index}`),
sourceSeenAt: null,
}))
const tombstoned = Array.from({ length: 2 }, (_, index) => ({
...row(`${connectorId}-z-tombstone-${index}`),
sourceSeenAt: null,
deletedAt: new Date(startedAt.getTime() - 60_000),
}))
await seed([...present, ...absent, ...tombstoned], present.length)
const { pass, stats, hardDeleted, walked } = await reconcile()
expect(pass).toMatchObject({ complete: true, holdNotice: null })
expect(stats.docsDeleted).toBe(absent.length)
expect(hardDeleted.sort()).toEqual(tombstoned.map((item) => item.id).sort())
const stored = await db
.select({ id: document.id, acl: document.acl, deletedAt: document.deletedAt })
.from(document)
.where(eq(document.connectorId, connectorId))
expect(stored).toHaveLength(present.length + absent.length)
for (const item of stored) expect(item.deletedAt !== null).toBe(item.id.includes('-z-'))
expect(
stored.filter((item) => item.id.includes('-z-')).every((item) => item.acl.length === 0)
).toBe(true)
await expectBounded(walked)
}, 120_000)

it('scans a dense window once and never reads every tombstone of the connector', async () => {
/**
* Three hard-delete pages of absent tombstones open the first window; the rest of the
* connector is tombstones the listing still sees, which the hard walk must pass over.
*/
const dense = Array.from({ length: 75 }, (_, index) => ({
...row(`${connectorId}-a-${String(index).padStart(3, '0')}`),
acl: [],
sourceSeenAt: null,
deletedAt: new Date(startedAt.getTime() - 60_000),
}))
const seenTombstones = Array.from({ length: 3 * RECONCILIATION_WINDOW_SIZE }, (_, index) => ({
...row(`${connectorId}-b-${String(index).padStart(6, '0')}`),
sourceSeenAt: startedAt,
deletedAt: new Date(startedAt.getTime() - 60_000),
}))
await seed([...dense, ...seenTombstones], seenTombstones.length)
const { pass, hardDeleted, walked } = await reconcile()
expect(pass).toMatchObject({ complete: true, holdNotice: null })
expect(hardDeleted.sort()).toEqual(dense.map((item) => item.id).sort())
expect(
await db
.select({ id: document.id })
.from(document)
.where(eq(document.connectorId, connectorId))
).toHaveLength(seenTombstones.length)
/**
* The only walk is the hard one: one statement per window, the dense first one included,
* then the tail, however many pages of matches a window holds.
*/
expect(walked).toHaveLength(
Math.floor((dense.length + seenTombstones.length) / RECONCILIATION_WINDOW_SIZE) + 1
)
/** `deleted_at < $1` implies the tombstone index, which would read every connector tombstone. */
await expectBounded(walked)
}, 120_000)
})

it('indexes a page, resumes under a new lease, and reconciles absence only after EOF', async () => {
await db
.update(knowledgeConnector)
Expand Down
63 changes: 32 additions & 31 deletions apps/sim/lib/knowledge/__integration__/scale.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,10 @@ import { readFileSync, statSync, writeFileSync } from 'node:fs'
import { db } from '@sim/db'
import { document, embedding, knowledgeConnector, user, workspace } from '@sim/db/schema'
import { createLogger } from '@sim/logger'
import { and, asc, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it } from 'vitest'
import { toStringOrNull } from '@sim/utils/coerce'
import { toArray, toRecord } from '@sim/utils/object'
import { and, eq, inArray, isNull, type SQL, sql } from 'drizzle-orm'
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
import { z } from 'zod'
import { resolveBillingAttribution } from '@/lib/billing/core/billing-attribution'
import {
Expand Down Expand Up @@ -139,6 +141,16 @@ async function explain(label: string, query: SQL, iterative = false) {
saveReport()
}

/** Every index a reported plan scans, at any depth. */
function planIndexNames(label: string): string[] {
const walk = (node: unknown): string[] => {
const record = toRecord(node)
const name = toStringOrNull(record['Index Name'])
return [...(name ? [name] : []), ...toArray(record.Plans).flatMap(walk)]
}
return toArray(report[`${label}.plan`]).flatMap((root) => walk(toRecord(root).Plan))
}

async function snapshot(label: string) {
const size =
await db.execute(sql`SELECT pg_database_size(current_database())::text AS database_bytes,
Expand Down Expand Up @@ -287,6 +299,8 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
await db.execute(
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}`
)
/** Plans the walks from statistics that see the rewritten absence, as autovacuum would. */
await db.execute(sql`ANALYZE document`)
await db
.update(knowledgeConnector)
.set({ listingCheckpoint: checkpoint })
Expand All @@ -309,35 +323,7 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
sql`SELECT id FROM document WHERE connector_id = ${ids.connectorId} AND user_excluded = false AND archived_at IS NULL
AND (source_seen_at IS NULL OR source_seen_at < ${startedAt.toISOString()}::timestamp) AND cardinality(acl) > 0 LIMIT 500`
)
const seenOrder = sql`COALESCE(${document.sourceSeenAt}, '-infinity'::timestamp)`
const absent = sql`connector_id = ${ids.connectorId} AND user_excluded = false AND archived_at IS NULL
AND ${seenOrder} < ${startedAt.toISOString()}::timestamp AND cardinality(acl) > 0`
const firstPage = db
.select({ id: document.id, seenAt: sql<string>`${seenOrder}::text` })
.from(document)
.where(absent)
.orderBy(asc(seenOrder), asc(document.id))
.limit(PAGE_SIZE)
await explain('reconciliation.keyset.first', firstPage.getSQL())
const firstCandidates = await firstPage
expect(firstCandidates).toHaveLength(PAGE_SIZE)
const nextPage = db
.select({ id: document.id })
.from(document)
.where(
and(
absent,
sql`(${seenOrder}, ${document.id}) > (${firstCandidates.at(-1)!.seenAt}::timestamp, ${firstCandidates.at(-1)!.id})`
)
)
.orderBy(asc(seenOrder), asc(document.id))
.limit(PAGE_SIZE)
await explain('reconciliation.keyset.next', nextPage.getSQL())
const nextCandidates = await nextPage
expect(nextCandidates).toHaveLength(PAGE_SIZE)
expect(new Set([...firstCandidates, ...nextCandidates].map((row) => row.id)).size).toBe(
2 * PAGE_SIZE
)
const statements = vi.spyOn(db.$client, 'unsafe')
const outcome = await measure('reconciliation.actual', async () =>
runConnectorContentPass({
connectorId: ids.connectorId,
Expand Down Expand Up @@ -368,8 +354,23 @@ describe.skipIf(!enabled)('knowledge scale: isolated real PostgreSQL, no provide
deadlineAt: Date.now() + 300_000,
})
)
/** Plans the first window scan the pass actually issued, so the plan follows the production SQL. */
const [windowScan] = statements.mock.calls.filter(([query]) =>
/^\s*with "document" as materialized/i.test(query)
)
statements.mockRestore()
expect(outcome.complete).toBe(true)
expect(outcome.holdNotice).toBeNull()
expect(windowScan).toBeDefined()
const [scanned] = await db.$client.unsafe(
`EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) ${windowScan[0]}`,
windowScan[1]
)
report['reconciliation.window.scan.plan'] = scanned['QUERY PLAN']
saveReport()
expect(planIndexNames('reconciliation.window.scan')).toContain(
'doc_connector_reconciliation_v2_idx'
)
const [removed] = await db.execute(
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`
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import {
sweepStaleMemberObservations,
writeProjectionPages,
} from '@/lib/knowledge/connectors/member-observations'
import { windowScan } from '@/lib/knowledge/connectors/reconciliation-window.test-helpers'
import { MEMBER_OBSERVATION_STALE_AFTER_HOURS } from '@/lib/knowledge/connectors/sync-limits'
import { type LeaseTransaction, SyncLockLostException } from '@/lib/knowledge/connectors/sync-lock'
import {
Expand Down Expand Up @@ -339,7 +340,8 @@ describe('applyMemberDocumentLifecycle', () => {
it('reports a reclaimed lease during a purge batch as the run being superseded', async () => {
dbChainMockFns.returning.mockResolvedValueOnce([])
queueTableRows(schemaMock.document, [])
queueTableRows(schemaMock.document, [])
/** The resurrection walk's only window, read through `db.execute`, is empty. */
dbChainMockFns.execute.mockResolvedValueOnce(windowScan([]))
queueTableRows(schemaMock.document, [{ id: 'd-1' }])
vi.mocked(hardDeleteDocuments).mockRejectedValueOnce(
new ConnectorSyncDeletionGuardError('lease reclaimed')
Expand Down
Loading
Loading