Skip to content

Commit 576408d

Browse files
authored
v0.8.55: projection index for sim search
2 parents a81b0e1 + 17e415d commit 576408d

6 files changed

Lines changed: 280 additions & 40 deletions

File tree

‎apps/sim/lib/knowledge/search/projection-source-acl-backfill.test.ts‎

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -78,12 +78,21 @@ describe('runProjectionSourceAclBackfill', () => {
7878

7979
it('analyzes and warms the projections on the same connection once both are filled, before closing it', async () => {
8080
await runProjectionSourceAclBackfill({})
81-
/** A row whose document is gone is not the fill's to finish; the probe joins the document. */
82-
expect(
83-
mockUnsafe.mock.calls.some(([query]) =>
84-
String(query).includes('JOIN document d ON d.id = s.document_id WHERE s.acl IS NULL')
81+
/**
82+
* A row whose document is gone is not the fill's to finish; the probe joins the document. It is
83+
* ordered and capped so only the unfilled-rows index can serve it: an `EXISTS` drops both and
84+
* leaves the planner a sequential scan of the projection.
85+
*/
86+
const probes = mockUnsafe.mock.calls
87+
.map(([query]) => String(query).replace(/\s+/g, ' '))
88+
.filter((query) => query.includes('AS unfilled'))
89+
expect(probes).toHaveLength(2)
90+
for (const probe of probes) {
91+
expect(probe).not.toContain('EXISTS')
92+
expect(probe).toContain(
93+
'JOIN document d ON d.id = s.document_id WHERE s.acl IS NULL ORDER BY s.id DESC LIMIT 1 ) IS NOT NULL AS unfilled'
8594
)
86-
).toBe(true)
95+
}
8796
expect(mockUnsafe.mock.calls.map(([query]) => query)).toEqual(
8897
expect.arrayContaining(['ANALYZE embedding_search', 'ANALYZE embedding_keyword_tin'])
8998
)

‎apps/sim/lib/knowledge/search/projection-source-acl-backfill.ts‎

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -148,14 +148,17 @@ export async function runProjectionSourceAclBackfill(
148148
/**
149149
* Whether no projection still holds a row the fill could give its source and ACL: a row without
150150
* them whose document exists. A row whose document is gone is not the fill's to finish and never
151-
* counts as left. Each read is one index probe while any such row remains.
151+
* counts as left. Each read is one probe of the unfilled-rows index while any such row remains:
152+
* ordered by id and capped at one row so the planner cannot take a sequential scan, which an
153+
* `EXISTS` would leave open by dropping the order and the limit.
152154
*/
153155
async function projectionsFilled(sql: postgres.Sql): Promise<boolean> {
154156
for (const projection of PROJECTION_SOURCE_ACL_TABLES) {
155157
const [row] = await sql.unsafe<Array<{ unfilled: boolean }>>(
156-
`SELECT EXISTS (
157-
SELECT 1 FROM ${projection} s JOIN document d ON d.id = s.document_id WHERE s.acl IS NULL
158-
) AS unfilled`
158+
`SELECT (
159+
SELECT s.id FROM ${projection} s JOIN document d ON d.id = s.document_id WHERE s.acl IS NULL
160+
ORDER BY s.id DESC LIMIT 1
161+
) IS NOT NULL AS unfilled`
159162
)
160163
if (row?.unfilled) return false
161164
}

‎apps/sim/lib/knowledge/search/queries.test.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2822,6 +2822,24 @@ describe('filters on a resolved scope', () => {
28222822
expect(caps.at(-1)).toBe('20000')
28232823
})
28242824

2825+
it('looks for an unfilled row through the ordered, capped read the partial index serves', async () => {
2826+
traversedRows = [{ id: 'a' }]
2827+
rerankRows = [hit('a', 'src-a')]
2828+
queueTableRows(schemaMock.embedding, rerankRows)
2829+
await handleVectorOnlySearch({
2830+
...params,
2831+
permitted: { kind: 'unbounded', broad: true },
2832+
accessPlan: plan(),
2833+
})
2834+
const probes = statements().filter((query) => query.sql.includes('AS unfilled'))
2835+
expect(probes).toHaveLength(1)
2836+
/** An `EXISTS` drops its order and limit, and the planner then takes a sequential scan. */
2837+
expect(probes[0].sql).not.toContain('EXISTS')
2838+
expect(probes[0].sql.replace(/\s+/g, ' ')).toContain(
2839+
'SELECT ( SELECT ? FROM ? WHERE ? IS NULL ORDER BY ? DESC LIMIT 1 ) IS NOT NULL AS unfilled'
2840+
)
2841+
})
2842+
28252843
it('tests the date through the document inside an on-row walk when the filtered set is unbounded', async () => {
28262844
traversedRows = [{ id: 'a' }]
28272845
rerankRows = [hit('a', 'src-a')]

‎apps/sim/lib/knowledge/search/queries.ts‎

Lines changed: 11 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -115,7 +115,13 @@ const PROJECTION_FILLED_TTL_MS = 60_000
115115

116116
/**
117117
* Whether the ranking projection still holds rows the backfill has not filled. Read off the
118-
* unfilled-rows index in microseconds and remembered briefly: the answer only ever changes once.
118+
* unfilled-rows index in milliseconds and remembered briefly: the answer only ever changes once.
119+
*
120+
* The read asks for the last unfilled row by id, not whether one exists: an `EXISTS` drops its
121+
* order and limit, and while most rows are unfilled the planner expects a sequential scan to
122+
* meet one at once, then walks the whole projection when the unfilled rows sit past the filled
123+
* ones. Ordered by id and capped at one row, the read can only be the partial index, whose
124+
* last entry is the row the fill reaches last.
119125
*/
120126
const projectionFilled = new LRUCache<
121127
ProjectionSourceAclTable,
@@ -134,7 +140,10 @@ const projectionFilled = new LRUCache<
134140
try {
135141
const [row] = await runSearchQuery(context.budget, context.stage, (executor) =>
136142
executor.execute<{ unfilled: boolean }>(sql`
137-
SELECT EXISTS (SELECT 1 FROM ${table} WHERE ${table.acl} IS NULL) AS unfilled`)
143+
SELECT (
144+
SELECT ${table.id} FROM ${table} WHERE ${table.acl} IS NULL
145+
ORDER BY ${table.id} DESC LIMIT 1
146+
) IS NOT NULL AS unfilled`)
138147
)
139148
return !row?.unfilled
140149
} catch {

‎packages/db/script-migrations/0021_embedding_search_connector.test.ts‎

Lines changed: 129 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,141 @@
11
/**
22
* @vitest-environment node
33
*/
4-
import { backfillProjectionSourceAcl } from '@sim/db/script-migrations/0021_embedding_search_connector'
4+
import {
5+
backfillProjectionSourceAcl,
6+
PROJECTION_SOURCE_ACL_PAGE_RETRIES,
7+
} from '@sim/db/script-migrations/0021_embedding_search_connector'
58
import type { Sql } from 'postgres'
6-
import { describe, expect, it, vi } from 'vitest'
9+
import { afterEach, describe, expect, it, vi } from 'vitest'
710

811
/** A session that must never be reached: every case below is refused before the first page. */
912
const untouched = { begin: vi.fn() } as unknown as Sql
1013

14+
type PageRow = { scanned: number; filled: number; last_id: string | null }
15+
16+
/** The message the database pairs with each cancellation SQLSTATE. */
17+
const CANCELLATION_MESSAGES: Record<string, string> = {
18+
'55P03': 'canceling statement due to lock timeout',
19+
'57014': 'canceling statement due to statement timeout',
20+
}
21+
22+
/** A driver error carrying a SQLSTATE, the shape `postgres` throws. */
23+
function postgresError(code: string, message = CANCELLATION_MESSAGES[code] ?? 'failed'): Error {
24+
return Object.assign(new Error(`${message} (SQLSTATE ${code})`), { code })
25+
}
26+
27+
/**
28+
* A session whose page statement answers from `outcomes` in order — a row, or an error to throw —
29+
* and records the cursor each page was bound to. `beforePage` runs with the page's index before
30+
* it answers.
31+
*/
32+
function sessionOf(outcomes: Array<PageRow | Error>, beforePage?: (index: number) => void) {
33+
const cursors: string[] = []
34+
const tx = {
35+
unsafe: vi.fn(async (query: string, params?: unknown[]) => {
36+
if (!query.includes('WITH page')) return []
37+
beforePage?.(cursors.length)
38+
cursors.push(String(params?.[0]))
39+
const outcome = outcomes.shift()
40+
if (outcome === undefined) throw new Error('No outcome left for this page')
41+
if (outcome instanceof Error) throw outcome
42+
return [outcome]
43+
}),
44+
}
45+
const session = {
46+
begin: vi.fn(async (work: (tx: unknown) => Promise<unknown>) => work(tx)),
47+
unsafe: vi.fn(async () => []),
48+
} as unknown as Sql
49+
return { session, cursors }
50+
}
51+
52+
/** Runs a backfill under fake timers, so its retry pauses pass without waiting. */
53+
async function backfillNow(...args: Parameters<typeof backfillProjectionSourceAcl>) {
54+
vi.useFakeTimers()
55+
const result = backfillProjectionSourceAcl(...args)
56+
/** A rejection must not surface as unhandled while the timers are still being drained. */
57+
const settled = result.catch(() => undefined)
58+
await vi.runAllTimersAsync()
59+
await settled
60+
return result
61+
}
62+
1163
describe('backfillProjectionSourceAcl', () => {
64+
afterEach(() => {
65+
vi.useRealTimers()
66+
})
67+
68+
it.each(['55P03', '57014'])(
69+
'retries the page after a %s timeout and moves the cursor only once it commits',
70+
async (code) => {
71+
const { session, cursors } = sessionOf([
72+
{ scanned: 2, filled: 2, last_id: 'id-2' },
73+
postgresError(code),
74+
postgresError(code),
75+
{ scanned: 1, filled: 1, last_id: 'id-3' },
76+
{ scanned: 0, filled: 0, last_id: null },
77+
])
78+
await expect(
79+
backfillNow(session, 'embedding_keyword_tin', { pauseMs: 0 })
80+
).resolves.toMatchObject({ scanned: 3, written: 3, afterId: 'id-3', done: true })
81+
expect(cursors).toEqual(['', 'id-2', 'id-2', 'id-2', 'id-3'])
82+
}
83+
)
84+
85+
it('gives up on a page that times out more than the retry limit in a row', async () => {
86+
const { session, cursors } = sessionOf(
87+
Array.from({ length: PROJECTION_SOURCE_ACL_PAGE_RETRIES + 1 }, () => postgresError('55P03'))
88+
)
89+
await expect(backfillNow(session, 'embedding_keyword_tin', { pauseMs: 0 })).rejects.toThrow(
90+
'SQLSTATE 55P03'
91+
)
92+
expect(cursors).toHaveLength(PROJECTION_SOURCE_ACL_PAGE_RETRIES + 1)
93+
expect(new Set(cursors)).toEqual(new Set(['']))
94+
})
95+
96+
it('propagates an error that is not a timeout without retrying', async () => {
97+
const { session, cursors } = sessionOf([postgresError('42P01')])
98+
await expect(backfillNow(session, 'embedding_search', { pauseMs: 0 })).rejects.toThrow(
99+
'SQLSTATE 42P01'
100+
)
101+
expect(cursors).toEqual([''])
102+
})
103+
104+
it('propagates an explicit cancellation, which shares the statement timeout SQLSTATE', async () => {
105+
const { session, cursors } = sessionOf([
106+
postgresError('57014', 'canceling statement due to user request'),
107+
])
108+
await expect(backfillNow(session, 'embedding_search', { pauseMs: 0 })).rejects.toThrow(
109+
'user request'
110+
)
111+
expect(cursors).toEqual([''])
112+
})
113+
114+
it('does not start another page when the budget ran out during the retry pause', async () => {
115+
const { session, cursors } = sessionOf([
116+
{ scanned: 1, filled: 1, last_id: 'id-1' },
117+
postgresError('57014'),
118+
])
119+
await expect(
120+
backfillNow(session, 'embedding_search', { pauseMs: 0, budgetMs: 1000 })
121+
).resolves.toMatchObject({ afterId: 'id-1', done: false })
122+
expect(cursors).toEqual(['', 'id-1'])
123+
})
124+
125+
it('leaves a page still failing at the budget to the continuation, from the last committed page', async () => {
126+
const { session, cursors } = sessionOf(
127+
[{ scanned: 1, filled: 1, last_id: 'id-1' }, postgresError('57014')],
128+
/** The second page spends the budget before the database cancels it. */
129+
(index) => {
130+
if (index === 1) vi.advanceTimersByTime(1000)
131+
}
132+
)
133+
await expect(
134+
backfillNow(session, 'embedding_search', { pauseMs: 0, budgetMs: 1000 })
135+
).resolves.toMatchObject({ afterId: 'id-1', done: false })
136+
expect(cursors).toEqual(['', 'id-1'])
137+
})
138+
12139
it.each([0, -1, 1.5, Number.NaN, Number.POSITIVE_INFINITY])(
13140
'refuses a page size of %s instead of reporting the projection filled',
14141
async (pageSize) => {

0 commit comments

Comments
 (0)