Skip to content

Commit acd0d95

Browse files
committed
improvement(knowledge): bound the projection warm and contain its extension probe
1 parent 5a099a7 commit acd0d95

4 files changed

Lines changed: 129 additions & 17 deletions

File tree

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

Lines changed: 73 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -95,6 +95,79 @@ describe('prewarmSearchProjection', () => {
9595
expect(warmed.map((item) => item.relation)).toEqual(['embedding_search_512_cosine_hnsw_idx'])
9696
})
9797

98+
it('returns nothing when the extension cannot be checked, never failing its caller', async () => {
99+
const fake = session({ installed: true })
100+
fake.unsafe = async () => {
101+
throw new Error('canceling statement due to user request')
102+
}
103+
await expect(prewarmSearchProjection(fake)).resolves.toEqual([])
104+
})
105+
106+
it('bounds every read by the budget left and leaves the rest cold once it is spent', async () => {
107+
vi.useFakeTimers()
108+
try {
109+
const fake = session({
110+
installed: true,
111+
relations: [
112+
'embedding_search',
113+
'embedding_keyword_tin',
114+
'embedding_search_512_cosine_hnsw_idx',
115+
],
116+
})
117+
const read = fake.unsafe
118+
fake.unsafe = async (query: string, parameters?: string[]) => {
119+
const rows = await read(query, parameters)
120+
/** Each read takes 400 ms of a 1 s budget. */
121+
if (query.includes('pg_prewarm(')) vi.advanceTimersByTime(400)
122+
return rows
123+
}
124+
const warmed = await prewarmSearchProjection(fake, { budgetMs: 1000 })
125+
expect(warmed.map((item) => item.relation)).toEqual([
126+
'embedding_search',
127+
'embedding_keyword_tin',
128+
'embedding_search_512_cosine_hnsw_idx',
129+
])
130+
const timeouts = fake.statements
131+
.filter((statement) => statement.query.startsWith('SET statement_timeout'))
132+
.map((statement) => Number(statement.query.split('= ')[1]))
133+
expect(timeouts).toEqual([1000, 600, 200])
134+
expect(fake.statements.at(-1)?.query).toBe('RESET statement_timeout')
135+
} finally {
136+
vi.useRealTimers()
137+
}
138+
})
139+
140+
it('skips the relations beyond a spent budget', async () => {
141+
vi.useFakeTimers()
142+
try {
143+
const fake = session({
144+
installed: true,
145+
relations: ['embedding_search', 'embedding_search_512_cosine_hnsw_idx'],
146+
})
147+
const read = fake.unsafe
148+
fake.unsafe = async (query: string, parameters?: string[]) => {
149+
const rows = await read(query, parameters)
150+
if (query.includes('pg_prewarm(')) vi.advanceTimersByTime(1500)
151+
return rows
152+
}
153+
const warmed = await prewarmSearchProjection(fake, { budgetMs: 1000 })
154+
expect(warmed.map((item) => item.relation)).toEqual(['embedding_search'])
155+
expect(
156+
fake.statements.filter((statement) => statement.query.includes('pg_prewarm('))
157+
).toHaveLength(1)
158+
} finally {
159+
vi.useRealTimers()
160+
}
161+
})
162+
163+
it('never sets a timeout on an unbounded pass', async () => {
164+
const fake = session({ installed: true, relations: ['embedding_search'] })
165+
await prewarmSearchProjection(fake)
166+
expect(fake.statements.some((statement) => statement.query.includes('statement_timeout'))).toBe(
167+
false
168+
)
169+
})
170+
98171
it('returns nothing when the catalog cannot be read, never failing its caller', async () => {
99172
const fake = session({ installed: true })
100173
fake.unsafe = async (query: string) => {

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

Lines changed: 46 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,15 @@ export interface PrewarmedRelation {
2222
elapsedMs: number
2323
}
2424

25+
export interface PrewarmOptions {
26+
/**
27+
* Wall-clock ceiling for the whole pass. Each read is bounded by the time left, and relations
28+
* beyond the ceiling stay cold; a caller with its own run limit sets it so warming can never
29+
* outlive the run that asked for it.
30+
*/
31+
budgetMs?: number
32+
}
33+
2534
/**
2635
* `pg_prewarm` is not a trusted extension, so the application role cannot create it and no
2736
* migration can; a superuser installs it once. Without it the projection warms only as searches
@@ -57,18 +66,23 @@ export async function prewarmRelation(
5766
* Heaps go first and the ranking indexes last, so where the cache cannot hold everything the
5867
* indexes are what survives: a walk reads far more index pages than heap pages. Relations are
5968
* resolved through the search path, so a schema that carries its own copy warms its own copy.
60-
* A relation that fails to warm is logged and skipped; warming is never worth failing the
69+
* Nothing here throws: a missing extension, an unreadable catalog, a relation that fails to
70+
* read or a spent budget is logged and skipped, since warming is never worth failing the
6171
* operation that asked for it.
6272
*/
6373
export async function prewarmSearchProjection(
64-
session: PrewarmSession
74+
session: PrewarmSession,
75+
options: PrewarmOptions = {}
6576
): Promise<PrewarmedRelation[]> {
66-
if (!(await pgPrewarmInstalled(session))) {
67-
logger.warn('pg_prewarm is not installed; the search projection warms only as it is searched')
68-
return []
69-
}
77+
const startedAt = Date.now()
78+
const remainingMs = () =>
79+
options.budgetMs === undefined ? undefined : options.budgetMs - (Date.now() - startedAt)
7080
let relations: string[]
7181
try {
82+
if (!(await pgPrewarmInstalled(session))) {
83+
logger.warn('pg_prewarm is not installed; the search projection warms only as it is searched')
84+
return []
85+
}
7286
relations = await rankingRelations(session)
7387
} catch (error) {
7488
logger.warn('Search projection relations could not be listed', {
@@ -77,20 +91,37 @@ export async function prewarmSearchProjection(
7791
return []
7892
}
7993
const warmed: PrewarmedRelation[] = []
80-
for (const relation of relations) {
81-
try {
82-
warmed.push(await prewarmRelation(session, relation))
83-
} catch (error) {
84-
logger.warn('Search projection relation failed to warm', {
85-
relation,
86-
error: getErrorMessage(error),
87-
})
94+
const cold: string[] = []
95+
try {
96+
for (const relation of relations) {
97+
const left = remainingMs()
98+
if (left !== undefined && left <= 0) {
99+
cold.push(relation)
100+
continue
101+
}
102+
try {
103+
if (left !== undefined) {
104+
await session.unsafe(`SET statement_timeout = ${Math.ceil(left)}`)
105+
}
106+
warmed.push(await prewarmRelation(session, relation))
107+
} catch (error) {
108+
cold.push(relation)
109+
logger.warn('Search projection relation failed to warm', {
110+
relation,
111+
error: getErrorMessage(error),
112+
})
113+
}
114+
}
115+
} finally {
116+
if (options.budgetMs !== undefined) {
117+
await Promise.resolve(session.unsafe('RESET statement_timeout')).catch(() => undefined)
88118
}
89119
}
90120
logger.info('Search projection warmed', {
91121
relations: warmed.length,
122+
cold,
92123
pages: warmed.reduce((sum, item) => sum + item.pages, 0),
93-
elapsedMs: warmed.reduce((sum, item) => sum + item.elapsedMs, 0),
124+
elapsedMs: Date.now() - startedAt,
94125
})
95126
return warmed
96127
}

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ vi.mock('@/lib/core/utils/background', () => ({
2828

2929
import {
3030
enqueueProjectionSourceAclBackfill,
31+
PROJECTION_PREWARM_BUDGET_MS,
3132
runProjectionSourceAclBackfill,
3233
} from '@/lib/knowledge/search/projection-source-acl-backfill'
3334

@@ -62,7 +63,7 @@ describe('runProjectionSourceAclBackfill', () => {
6263
it('warms the projections on the same connection once both are filled, before closing it', async () => {
6364
await runProjectionSourceAclBackfill({})
6465
expect(mockPrewarm).toHaveBeenCalledTimes(1)
65-
expect(mockPrewarm).toHaveBeenCalledWith(connection)
66+
expect(mockPrewarm).toHaveBeenCalledWith(connection, { budgetMs: PROJECTION_PREWARM_BUDGET_MS })
6667
expect(mockPrewarm.mock.invocationCallOrder[0]).toBeLessThan(
6768
mockEnd.mock.invocationCallOrder[0]
6869
)

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

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,13 @@ const logger = createLogger('ProjectionSourceAclBackfill')
1616

1717
export const PROJECTION_SOURCE_ACL_BACKFILL_TASK_ID = 'projection-source-acl-backfill'
1818

19+
/**
20+
* Ceiling on warming the projections after the fill. It sits inside the headroom the worker
21+
* keeps beyond a run's fill budget, so a slow read of a large projection can never carry the
22+
* completed run past the worker's limit and repeat the fill on retry.
23+
*/
24+
export const PROJECTION_PREWARM_BUDGET_MS = 15 * 60 * 1000
25+
1926
/** Where a run stopped, so the next one carries on from there instead of rescanning. */
2027
export interface ProjectionSourceAclBackfillCursor {
2128
projection: ProjectionSourceAclTable
@@ -71,7 +78,7 @@ export async function runProjectionSourceAclBackfill(
7178
elapsedMs: Date.now() - startedAt,
7279
})
7380
/** The fill just streamed through both projections; put the ranking pages back before anyone searches. */
74-
await prewarmSearchProjection(sql)
81+
await prewarmSearchProjection(sql, { budgetMs: PROJECTION_PREWARM_BUDGET_MS })
7582
return null
7683
} finally {
7784
await sql.end()

0 commit comments

Comments
 (0)