Skip to content

Commit 2b77b50

Browse files
committed
fix(file-search): bound chunk inserts by trigram work
1 parent eaa9ff0 commit 2b77b50

7 files changed

Lines changed: 242 additions & 41 deletions

File tree

‎apps/sim/lib/workspace-files/search/README.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,9 @@ PostgreSQL stores complete extracted text in bounded chunks. Object storage rema
1010

1111
8 KiB values may still use PostgreSQL TOAST. The bound controls the size of each logical value and detoast operation; avoiding TOAST entirely is not the objective. Tiny lines share rows, so row count scales with bytes instead of newline count. Worst-case line packing can leave roughly half a block unused; long-line overlap adds at most eight bytes per fragment.
1212

13-
Workers download and extract outside database transactions, then insert batches of at most 250 rows / 128 KiB. Each batch checks the build token and lease. Publication locks the canonical file, build, and revision in that order, verifies the stored chunk count, and changes the visible pointer only after every batch succeeds. Old dispatch failure callbacks cannot overwrite newer dispatches or successful builds.
13+
Workers download and extract outside database transactions, then insert batches of at most 250 rows / 128 KiB and a target of 8,192 estimated trigram keys. Each batch checks the build token and lease. Publication locks the canonical file, build, and revision in that order, verifies the stored chunk count, and changes the visible pointer only after every batch succeeds. Old dispatch failure callbacks cannot overwrite newer dispatches or successful builds.
1414

15-
The chunk GIN index uses `fastupdate = off`. Each bounded insert updates the main index directly instead of appending to a shared pending list. With deferred updates enabled, even a small insert can cross the pending-list threshold and synchronously merge accumulated work from other files. Direct updates trade some bulk-write throughput for avoiding that foreground cleanup cliff. They do not eliminate normal index I/O, vacuum, or storage contention; the row and worker limits still apply. The 128 KiB batch budget bounds direct index work per transaction without changing the 25 MiB file coverage limit. Dense text with many distinct trigrams and an index working set larger than the available cache can still exceed the statement deadline. Capacity validation must include that cache pressure, not only a small corpus or a row count. A batch that takes at least two seconds, including one a statement timeout cancels, is logged with its rows, bytes, and estimated trigram key count (pg_trgm's extraction, reproduced exactly for the database's `en_US.UTF-8` ctype), so a key-based batch budget can be sized from production timings. A failed query is reported to the task runner by error code only; Drizzle's message carries the bound file text.
15+
The chunk GIN index uses `fastupdate = off`. Each bounded insert updates the main index directly instead of appending to a shared pending list. With deferred updates enabled, even a small insert can cross the pending-list threshold and synchronously merge accumulated work from other files. Direct updates trade some bulk-write throughput for avoiding that foreground cleanup cliff. They do not eliminate normal index I/O, vacuum, or storage contention; the row and worker limits still apply. Bytes alone do not bound GIN posting updates: dense text can generate thousands of distinct keys in each chunk. The batch planner sums each chunk's distinct trigram count, including repeated keys across rows, and flushes before the next chunk exceeds the key target. A single 8 KiB chunk is always allowed to make progress even if word padding pushes its estimate slightly above the target; it is written alone. Storage validates the same budget before opening the transaction. These are write scheduling bounds, not file exclusions: chunk boundaries, complete-file publication, and the 25 MiB file coverage limit are unchanged. The estimate mirrors pg_trgm under `en_US.UTF-8`; it is a work estimate, not a latency guarantee. Capacity validation must still include dense text, concurrent writers, and an index working set larger than available cache. A batch that takes at least two seconds, including a canceled statement, is logged with its rows, bytes, and estimated key count. A failed query is reported to the task runner by error code only; Drizzle's message carries the bound file text.
1616

1717
Indexing transactions have separate limits from search: ten seconds per statement, five seconds waiting for a lock, and thirty seconds total on PostgreSQL 17. The outer limit leaves time for ordinary statement cancellation and rollback instead of terminating the connection at the same ten-second deadline. PostgreSQL 16 uses the compatible idle-transaction guard. A canceled batch remains unpublished; the existing task retry starts a fresh fenced build, and cleanup retires the previous attempt. This does not automatically retry revisions already marked failed.
1818

‎apps/sim/lib/workspace-files/search/chunks.integration.ts‎

Lines changed: 41 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -32,16 +32,15 @@ vi.mock('@/lib/copilot/tools/server/files/doc-compile', () => ({ resolveServable
3232
vi.mock('@/lib/file-parsers', () => ({ parseBuffer: vi.fn(), isSupportedFileType: vi.fn() }))
3333

3434
import {
35-
FILE_SEARCH_CHUNK_BYTES,
3635
FILE_SEARCH_CLEANUP_BATCH_ROWS,
3736
FILE_SEARCH_CLEANUP_BUDGET_MS,
3837
FILE_SEARCH_CLEANUP_MAX_BATCHES,
39-
FILE_SEARCH_INSERT_BATCH_BYTES,
40-
FILE_SEARCH_INSERT_BATCH_ROWS,
38+
FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS,
4139
FILE_SEARCH_QUERY_GLOBAL_CONCURRENCY,
4240
FILE_SEARCH_QUERY_WORKSPACE_CONCURRENCY,
4341
} from '@/lib/workspace-files/search/constants'
4442
import { prepareWorkspaceFileSearchDispatch } from '@/lib/workspace-files/search/dispatcher'
43+
import { iterateFileSearchBatches } from '@/lib/workspace-files/search/index-batches'
4544
import {
4645
iterateFileSearchChunks,
4746
planFileSearchIndex,
@@ -142,14 +141,8 @@ describe('chunked workspace file search on PostgreSQL', () => {
142141
expect(build).not.toBeNull()
143142
const plan = planFileSearchIndex({ text, partial: false }, signal)
144143
const chunks = [...iterateFileSearchChunks(plan, signal)]
145-
const batchRows = Math.min(
146-
FILE_SEARCH_INSERT_BATCH_ROWS,
147-
Math.floor(FILE_SEARCH_INSERT_BATCH_BYTES / FILE_SEARCH_CHUNK_BYTES)
148-
)
149-
for (let offset = 0; offset < chunks.length; offset += batchRows)
150-
expect(
151-
await appendFileSearchChunks(build!, chunks.slice(offset, offset + batchRows), signal)
152-
).toBe(true)
144+
for (const batch of iterateFileSearchBatches(chunks, signal))
145+
expect(await appendFileSearchChunks(build!, batch, signal)).toBe(true)
153146
expect(
154147
await publishFileSearchBuild(
155148
build!,
@@ -304,6 +297,43 @@ describe('chunked workspace file search on PostgreSQL', () => {
304297
expect((await search('needle')).results).toMatchObject([{ fileId: 'file-1', lineNumber: 1 }])
305298
})
306299

300+
it('bounds native GIN posting work and keeps dense-file exact and regex line results', async () => {
301+
const lines = Array.from(
302+
{ length: 1200 },
303+
(_, i) => `dependency-${i}: sha512-${createHash('sha512').update(String(i)).digest('base64')}`
304+
)
305+
const text = lines.join('\n')
306+
const plan = planFileSearchIndex({ text, partial: false }, signal)
307+
const chunks = [...iterateFileSearchChunks(plan, signal)]
308+
const build = (await beginFileSearchBuild(revision))!
309+
await expect(appendFileSearchChunks(build, chunks.slice(0, 16), signal)).rejects.toThrow(
310+
'insert batch exceeds its budget'
311+
)
312+
for (const batch of iterateFileSearchBatches(chunks, signal)) {
313+
const [{ keys }] = await connection`SELECT sum(cardinality(show_trgm(content)))::int AS keys
314+
FROM (VALUES ${connection(batch.map((c) => [c.content]))}) AS chunk(content)`
315+
expect(keys).toBeLessThanOrEqual(FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS)
316+
expect(await appendFileSearchChunks(build, batch, signal)).toBe(true)
317+
}
318+
expect((await search('dependency-1199')).results).toEqual([])
319+
expect(
320+
await publishFileSearchBuild(
321+
build,
322+
{
323+
status: 'ready',
324+
chunkCount: chunks.length,
325+
lineCount: plan.lineCount,
326+
indexedBytes: plan.indexedBytes,
327+
},
328+
signal
329+
)
330+
).toBe(true)
331+
expect((await search(lines[1199])).results).toMatchObject([{ lineNumber: 1200 }])
332+
expect(
333+
(await search('^dependency-1199: sha512-[A-Za-z0-9+/]+=*$', 'regex')).results
334+
).toMatchObject([{ lineNumber: 1200 }])
335+
})
336+
307337
it('packs a million short lines without a million rows and bounds every stored value', async () => {
308338
await index('abc\n'.repeat(1_000_000))
309339
const [row] =

‎apps/sim/lib/workspace-files/search/constants.ts‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -61,8 +61,13 @@ export const FILE_SEARCH_CLEANUP_MIN_BATCH_MS =
6161
FILE_SEARCH_CLEANUP_BUDGET_MS / FILE_SEARCH_CLEANUP_MAX_BATCHES
6262
export const FILE_SEARCH_RECONCILE_INTERVAL_MS = 60 * 60 * 1000
6363
export const FILE_SEARCH_INSERT_BATCH_ROWS = 250
64-
/** Direct GIN writes perform index work in each insert, so transactions use smaller byte batches. */
64+
/** Bounds the text payload independently of its index work. */
6565
export const FILE_SEARCH_INSERT_BATCH_BYTES = 128 * 1024
66+
/**
67+
* Sum of each chunk's distinct trigram keys, bounding direct GIN posting updates per insert.
68+
* A single 8 KiB chunk may exceed this target slightly due to word padding and is written alone.
69+
*/
70+
export const FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS = 8 * 1024
6671
/** Batches at least this slow are logged with their estimated trigram key count. */
6772
export const FILE_SEARCH_SLOW_INSERT_BATCH_MS = 2000
6873

Lines changed: 126 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,126 @@
1+
/** @vitest-environment node */
2+
import { createHash } from 'node:crypto'
3+
import { describe, expect, it } from 'vitest'
4+
import {
5+
FILE_SEARCH_CHUNK_BYTES,
6+
FILE_SEARCH_INSERT_BATCH_BYTES,
7+
FILE_SEARCH_INSERT_BATCH_ROWS,
8+
FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS,
9+
} from '@/lib/workspace-files/search/constants'
10+
import { iterateFileSearchBatches } from '@/lib/workspace-files/search/index-batches'
11+
import {
12+
estimateTrigramKeys,
13+
type FileSearchChunk,
14+
iterateFileSearchChunks,
15+
planFileSearchIndex,
16+
} from '@/lib/workspace-files/search/index-plan'
17+
18+
const signal = new AbortController().signal
19+
const chunk = (content: string, ordinal = 0): FileSearchChunk => ({
20+
content,
21+
ordinal,
22+
lineStart: ordinal + 1,
23+
fragment: false,
24+
overlap: 0,
25+
})
26+
27+
/** A de Bruijn prefix has no repeated three-letter windows, exercising the padding overhead. */
28+
function uniqueTrigrams(): string {
29+
const alphabet = 'abcdefghijklmnopqrstuvwxyz'
30+
const positions = Array<number>(4).fill(0)
31+
let text = ''
32+
const visit = (depth: number, period: number) => {
33+
if (depth > 3) {
34+
if (3 % period === 0) for (let i = 1; i <= period; i++) text += alphabet[positions[i]]
35+
return
36+
}
37+
positions[depth] = positions[depth - period]
38+
visit(depth + 1, period)
39+
for (let i = positions[depth - period] + 1; i < alphabet.length; i++) {
40+
positions[depth] = i
41+
visit(depth + 1, depth)
42+
}
43+
}
44+
visit(1, 1)
45+
return text.slice(0, FILE_SEARCH_CHUNK_BYTES)
46+
}
47+
48+
describe('file search insert batches', () => {
49+
it('counts GIN work per row, preserving every dense lockfile line and ordinal', () => {
50+
const text = Array.from(
51+
{ length: 1200 },
52+
(_, i) =>
53+
`dependency-${i}: sha512-${createHash('sha512').update(String(i)).digest('base64')}\n`
54+
).join('')
55+
const chunks = [
56+
...iterateFileSearchChunks(planFileSearchIndex({ text, partial: false }, signal), signal),
57+
]
58+
const batches = [...iterateFileSearchBatches(chunks, signal)]
59+
expect(batches.length).toBeGreaterThan(
60+
Math.ceil(Buffer.byteLength(text) / FILE_SEARCH_INSERT_BATCH_BYTES)
61+
)
62+
expect(batches.flat()).toEqual(chunks)
63+
expect(
64+
batches
65+
.flat()
66+
.map((c) => c.content)
67+
.join('')
68+
).toBe(text)
69+
for (const batch of batches) {
70+
expect(batch.reduce((sum, c) => sum + estimateTrigramKeys(c.content), 0)).toBeLessThanOrEqual(
71+
FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS
72+
)
73+
}
74+
const repeated = chunk(chunks[0].content)
75+
const sameContent = [...iterateFileSearchBatches(Array(16).fill(repeated), signal)]
76+
expect(sameContent.length).toBeGreaterThan(1)
77+
})
78+
79+
it('retains byte batching for low-key text and the row bound for tiny chunks', () => {
80+
const full = chunk('a'.repeat(FILE_SEARCH_CHUNK_BYTES))
81+
const byBytes = [...iterateFileSearchBatches(Array(17).fill(full), signal)]
82+
expect(byBytes.map((batch) => batch.length)).toEqual([16, 1])
83+
const byRows = [
84+
...iterateFileSearchBatches(Array(FILE_SEARCH_INSERT_BATCH_ROWS + 1).fill(chunk('')), signal),
85+
]
86+
expect(byRows.map((batch) => batch.length)).toEqual([FILE_SEARCH_INSERT_BATCH_ROWS, 1])
87+
})
88+
89+
it('writes a maximum-key chunk alone without splitting or dropping it', () => {
90+
const dense = chunk(uniqueTrigrams(), 1)
91+
expect(estimateTrigramKeys(dense.content)).toBeGreaterThan(
92+
FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS
93+
)
94+
const first = chunk('first')
95+
const last = chunk('last', 2)
96+
expect([...iterateFileSearchBatches([first, dense, last], signal)]).toEqual([
97+
[first],
98+
[dense],
99+
[last],
100+
])
101+
})
102+
103+
it('rejects oversized chunks and emits no empty batch', () => {
104+
expect([...iterateFileSearchBatches([], signal)]).toEqual([])
105+
expect(() => [
106+
...iterateFileSearchBatches([chunk('a'.repeat(FILE_SEARCH_CHUNK_BYTES + 1))], signal),
107+
]).toThrow('chunk exceeds its budget')
108+
})
109+
110+
it('consumes only one batch plus lookahead and honors cancellation between writes', () => {
111+
let consumed = 0
112+
function* source() {
113+
for (let i = 0; i < 100; i++) {
114+
consumed++
115+
yield chunk('a'.repeat(FILE_SEARCH_CHUNK_BYTES), i)
116+
}
117+
}
118+
const controller = new AbortController()
119+
const batches = iterateFileSearchBatches(source(), controller.signal)
120+
expect(batches.next().value).toHaveLength(16)
121+
expect(consumed).toBe(17)
122+
controller.abort(new Error('canceled'))
123+
expect(() => batches.next()).toThrow('canceled')
124+
expect(consumed).toBe(17)
125+
})
126+
})
Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
import { Buffer } from 'node:buffer'
2+
import {
3+
FILE_SEARCH_CHUNK_BYTES,
4+
FILE_SEARCH_INSERT_BATCH_BYTES,
5+
FILE_SEARCH_INSERT_BATCH_ROWS,
6+
FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS,
7+
} from '@/lib/workspace-files/search/constants'
8+
import { estimateTrigramKeys, type FileSearchChunk } from '@/lib/workspace-files/search/index-plan'
9+
10+
/** A single bounded chunk must make progress even when its estimate exceeds the key target. */
11+
export function exceedsFileSearchBatchBudget(rows: number, bytes: number, keys: number): boolean {
12+
return (
13+
rows > FILE_SEARCH_INSERT_BATCH_ROWS ||
14+
bytes > FILE_SEARCH_INSERT_BATCH_BYTES ||
15+
(rows > 1 && keys > FILE_SEARCH_INSERT_BATCH_TRIGRAM_KEYS)
16+
)
17+
}
18+
19+
/**
20+
* Bounds each insert's payload and estimated GIN work without changing stored chunks or coverage.
21+
* Keys are counted per row: repeated keys across rows still require separate posting updates.
22+
*/
23+
export function* iterateFileSearchBatches(
24+
chunks: Iterable<FileSearchChunk>,
25+
signal: AbortSignal
26+
): Generator<FileSearchChunk[]> {
27+
let batch: FileSearchChunk[] = []
28+
let bytes = 0
29+
let keys = 0
30+
for (const chunk of chunks) {
31+
signal.throwIfAborted()
32+
const chunkBytes = Buffer.byteLength(chunk.content)
33+
if (chunkBytes > FILE_SEARCH_CHUNK_BYTES)
34+
throw new Error('File search chunk exceeds its budget')
35+
const chunkKeys = estimateTrigramKeys(chunk.content)
36+
if (
37+
batch.length &&
38+
exceedsFileSearchBatchBudget(batch.length + 1, bytes + chunkBytes, keys + chunkKeys)
39+
) {
40+
yield batch
41+
signal.throwIfAborted()
42+
batch = []
43+
bytes = keys = 0
44+
}
45+
batch.push(chunk)
46+
bytes += chunkBytes
47+
keys += chunkKeys
48+
}
49+
signal.throwIfAborted()
50+
if (batch.length) yield batch
51+
}

‎apps/sim/lib/workspace-files/search/index-state.ts‎

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -10,16 +10,16 @@ import { and, eq, sql } from 'drizzle-orm'
1010
import type { DbTransaction } from '@/lib/db/types'
1111
import {
1212
FILE_SEARCH_BUILD_LEASE_MS,
13+
FILE_SEARCH_CHUNK_BYTES,
1314
FILE_SEARCH_CLEANUP_BATCH_BUILDS,
1415
FILE_SEARCH_CLEANUP_BATCH_ROWS,
1516
FILE_SEARCH_CLEANUP_BUDGET_MS,
1617
FILE_SEARCH_CLEANUP_MAX_BATCHES,
1718
FILE_SEARCH_CLEANUP_MIN_BATCH_MS,
1819
FILE_SEARCH_INDEX_TRANSACTION_LIMITS,
19-
FILE_SEARCH_INSERT_BATCH_BYTES,
20-
FILE_SEARCH_INSERT_BATCH_ROWS,
2120
} from '@/lib/workspace-files/search/constants'
22-
import type { FileSearchChunk } from '@/lib/workspace-files/search/index-plan'
21+
import { exceedsFileSearchBatchBudget } from '@/lib/workspace-files/search/index-batches'
22+
import { estimateTrigramKeys, type FileSearchChunk } from '@/lib/workspace-files/search/index-plan'
2323
import { configureFileSearchTransaction } from '@/lib/workspace-files/search/transaction'
2424

2525
export interface FileSearchRevision {
@@ -139,7 +139,7 @@ export async function beginFileSearchBuild(
139139
})
140140
}
141141

142-
/** Each batch is fenced and byte-bounded; no file or parser work runs inside this transaction. */
142+
/** Each batch is fenced and work-bounded; no file or parser work runs inside this transaction. */
143143
export async function appendFileSearchChunks(
144144
build: FileSearchBuild,
145145
chunks: readonly FileSearchChunk[],
@@ -148,9 +148,12 @@ export async function appendFileSearchChunks(
148148
signal.throwIfAborted()
149149
if (!chunks.length) return true
150150
if (
151-
chunks.length > FILE_SEARCH_INSERT_BATCH_ROWS ||
152-
chunks.reduce((sum, c) => sum + Buffer.byteLength(c.content), 0) >
153-
FILE_SEARCH_INSERT_BATCH_BYTES
151+
chunks.some((chunk) => Buffer.byteLength(chunk.content) > FILE_SEARCH_CHUNK_BYTES) ||
152+
exceedsFileSearchBatchBudget(
153+
chunks.length,
154+
chunks.reduce((sum, chunk) => sum + Buffer.byteLength(chunk.content), 0),
155+
chunks.reduce((sum, chunk) => sum + estimateTrigramKeys(chunk.content), 0)
156+
)
154157
) {
155158
throw new Error('File search insert batch exceeds its budget')
156159
}

‎apps/sim/lib/workspace-files/search/indexing.ts‎

Lines changed: 6 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -5,12 +5,11 @@ import { redactDatabaseQueryError } from '@/lib/core/errors/database-query-error
55
import { isPayloadSizeLimitError } from '@/lib/core/utils/stream-limits'
66
import { getWorkspaceFile } from '@/lib/uploads/contexts/workspace'
77
import {
8-
FILE_SEARCH_INSERT_BATCH_BYTES,
9-
FILE_SEARCH_INSERT_BATCH_ROWS,
108
FILE_SEARCH_MAX_SOURCE_BYTES,
119
FILE_SEARCH_SLOW_INSERT_BATCH_MS,
1210
} from '@/lib/workspace-files/search/constants'
1311
import { extractIndexText, loadIndexableBytes } from '@/lib/workspace-files/search/extract'
12+
import { iterateFileSearchBatches } from '@/lib/workspace-files/search/index-batches'
1413
import {
1514
estimateTrigramKeys,
1615
type FileSearchChunk,
@@ -46,7 +45,7 @@ function parseRevision(payload: WorkspaceFileSearchIndexPayload): FileSearchRevi
4645

4746
/**
4847
* Appends one batch and records slow ones, including a batch a statement timeout cancels, with the
49-
* trigram key load that drives direct GIN insert cost. Keys are estimated only for slow batches.
48+
* trigram key load that drives direct GIN insert cost. Logging never includes the indexed text.
5049
*/
5150
async function appendTimedBatch(
5251
build: FileSearchBuild,
@@ -107,25 +106,12 @@ export async function indexWorkspaceFileForSearch(
107106
return
108107
}
109108
const plan = planFileSearchIndex(extracted, signal)
110-
let batch: FileSearchChunk[] = []
111-
let batchBytes = 0
112109
let chunkCount = 0
113-
for (const chunk of iterateFileSearchChunks(plan, signal)) {
114-
const chunkBytes = Buffer.byteLength(chunk.content, 'utf8')
115-
if (
116-
batch.length &&
117-
(batch.length >= FILE_SEARCH_INSERT_BATCH_ROWS ||
118-
batchBytes + chunkBytes > FILE_SEARCH_INSERT_BATCH_BYTES)
119-
) {
120-
if (!(await appendTimedBatch(build, batch, batchBytes, signal))) return
121-
batch = []
122-
batchBytes = 0
123-
}
124-
batch.push(chunk)
125-
batchBytes += chunkBytes
126-
chunkCount++
110+
for (const batch of iterateFileSearchBatches(iterateFileSearchChunks(plan, signal), signal)) {
111+
const batchBytes = batch.reduce((sum, chunk) => sum + Buffer.byteLength(chunk.content), 0)
112+
if (!(await appendTimedBatch(build, batch, batchBytes, signal))) return
113+
chunkCount += batch.length
127114
}
128-
if (!(await appendTimedBatch(build, batch, batchBytes, signal))) return
129115
const published = await publishFileSearchBuild(
130116
build,
131117
{ status: 'ready', chunkCount, lineCount: plan.lineCount, indexedBytes: plan.indexedBytes },

0 commit comments

Comments
 (0)