Skip to content

Commit 1f16b03

Browse files
fix(file-search): avoid foreground GIN cleanup stalls (#7995)
* fix(file-search): avoid foreground GIN cleanup stalls * fix(file-search): create chunk GIN index concurrently
1 parent 99d6a9f commit 1f16b03

8 files changed

Lines changed: 27942 additions & 24 deletions

File tree

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

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,11 @@ 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 / 1 MiB. 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. 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.
14+
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.
16+
17+
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.
1418

1519
The indexing task uses an isolated `medium-2x` Trigger worker (4 GB RAM). Document parsers can materialize expanded content before chunking, so source and extracted-text byte limits do not bound parser memory. Parser complexity guards and the worker's memory budget remain separate protections.
1620

@@ -41,16 +45,22 @@ Search results retain `fileId`, 1-based `lineNumber`, and bounded `text` preview
4145
3. Before retiring legacy storage, verify the new app and Trigger workers are fully deployed, old runs/retries have drained, the backfill cursor has completed, and scoped coverage is ready or explicitly excluded. Investigate failed or stale pending revisions. Check cleanup backlog and run representative exact/regex searches, including long lines and folder scopes.
4246
4. After the rollback window, ship a separate contract PR removing the legacy schema and dropping `workspace_file_search_segment` / `workspace_file_search_index` with a short lock timeout. Do not delete the entire old index row-by-row or backfill it inside the schema migration. Dropping obsolete tables reclaims their heap, indexes, and TOAST together. The `contract-pending` marker in `packages/db/schema.ts` tracks this step.
4347

44-
Until that contract deploy, legacy foreign-key cascades can still make a hard file/workspace deletion expensive. New-index cleanup is bounded; retaining the old schema cannot erase that legacy cost. The earlier timestamp-repair script detects the chunk schema and leaves obsolete legacy text for this contract step instead of deleting it in bulk. No production cleanup is part of this PR.
48+
Until that contract deploy, legacy foreign-key cascades can still make a hard file/workspace deletion expensive. New-index cleanup is bounded; retaining the old schema cannot erase that legacy cost. The earlier timestamp-repair script detects the chunk schema and leaves obsolete legacy text for this contract step instead of deleting it in bulk. Legacy-table retirement remains a separate contract migration.
4549

4650
Rollback before retirement requires restoring the old trigger function as well as the old app/worker version, and reconciling legacy revisions written during the cutover. Do not assume retained tables are automatically up to date. Canonical revision joins prevent stale content from being returned.
4751

52+
### Direct GIN writes
53+
54+
Migration `0364_workspace_file_search_direct_gin.sql` changes the existing index's storage option without rebuilding it or rewriting chunks. It commits that change before calling PostgreSQL's `gin_clean_pending_list` with a five-minute statement budget. Turning off deferred updates alone leaves the previous pending list in place, so this one-time drain is required. Searches continue to see both existing pending entries and main-index entries throughout the transition, and old workers use the same schema. Deploy the smaller-batch worker policy in the same release; running older workers retain their previous batch size until replaced.
55+
56+
If maintenance times out, the storage option remains off and migration replay safely resumes the drain; no content is discarded and no new pending entries accumulate. The migration uses the runner's direct connection and short DDL lock timeout. It adds no maintenance scheduler. After deployment, check that the index is valid, `reloptions` includes `fastupdate=off`, the migration completed, and chunk insert latency and timeout rates remain healthy under the existing worker concurrency. Previously failed revisions require an explicitly scoped retry after their failure reason and current content version are checked; no blanket backfill or deletion is part of this change.
57+
4858
## Verification
4959

50-
Run unit tests in `apps/sim` with `bunx vitest run lib/workspace-files/search lib/file-parsers`. Run the PostgreSQL suites on both PostgreSQL 16 and 17 against a disposable local database through `KNOWLEDGE_ACL_TEST_DATABASE_URL` and `--mode integration`. `chunks.integration.ts` applies the actual trigger migrations in an isolated schema. It covers build fencing, revision changes, deletion, cleanup bounds, complete-line matching, UTF-8 boundaries, scope, and admission limits.
60+
Run unit tests in `apps/sim` with `bunx vitest run lib/workspace-files/search lib/file-parsers`. Run the PostgreSQL suites on both PostgreSQL 16 and 17 against a disposable local database through `KNOWLEDGE_ACL_TEST_DATABASE_URL` and `--mode integration`. `chunks.integration.ts` applies the actual trigger and index migrations in an isolated schema. It covers build fencing, revision changes, deletion, cleanup bounds, complete-line matching, UTF-8 boundaries, scope, admission limits, GIN migration replay with a populated pending list, and statement cancellation followed by a successful retry on an intact connection. Its local database role must be able to install the `pgstattuple` diagnostic extension.
5161

5262
Set `FILE_SEARCH_BENCHMARK_FILES` to change the synthetic file count (default 1,000, maximum 10,000). Set `FILE_SEARCH_BENCHMARK_OUTPUT` to an output path when running the chunk integration suite to record repeated end-to-end searches and `EXPLAIN (ANALYZE, BUFFERS)` plans on a synthetic multi-file corpus. The fixture is synthetic; it contains no production content.
5363

5464
## PostgreSQL references
5565

56-
The design uses PostgreSQL's documented [TOAST behavior](https://www.postgresql.org/docs/17/storage-toast.html), [trigram index support for LIKE and regex](https://www.postgresql.org/docs/17/pgtrgm.html), and [EXPLAIN guidance](https://www.postgresql.org/docs/17/using-explain.html). The chunk size and query paths are application choices validated by the synthetic fixture, not PostgreSQL hard limits.
66+
The design uses PostgreSQL's documented [TOAST behavior](https://www.postgresql.org/docs/17/storage-toast.html), [trigram index support for LIKE and regex](https://www.postgresql.org/docs/17/pgtrgm.html), [GIN pending-list tradeoffs](https://www.postgresql.org/docs/17/gin.html#GIN-FAST-UPDATE), [index storage parameters](https://www.postgresql.org/docs/17/sql-createindex.html), and [EXPLAIN guidance](https://www.postgresql.org/docs/17/using-explain.html). The chunk size and query paths are application choices validated by the synthetic fixture, not PostgreSQL hard limits.

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

Lines changed: 125 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -32,9 +32,12 @@ 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,
3536
FILE_SEARCH_CLEANUP_BATCH_ROWS,
3637
FILE_SEARCH_CLEANUP_BUDGET_MS,
3738
FILE_SEARCH_CLEANUP_MAX_BATCHES,
39+
FILE_SEARCH_INSERT_BATCH_BYTES,
40+
FILE_SEARCH_INSERT_BATCH_ROWS,
3841
FILE_SEARCH_QUERY_GLOBAL_CONCURRENCY,
3942
FILE_SEARCH_QUERY_WORKSPACE_CONCURRENCY,
4043
} from '@/lib/workspace-files/search/constants'
@@ -93,6 +96,35 @@ describe('chunked workspace file search on PostgreSQL', () => {
9396
onnotice: () => {},
9497
})
9598
)
99+
const ginWriteMigration = '0364_workspace_file_search_direct_gin.sql'
100+
101+
async function applyMigration(migration: string) {
102+
const source = readFileSync(
103+
resolve(process.cwd(), '../../packages/db/migrations', migration),
104+
'utf8'
105+
).replaceAll('"public".', `"${schema}".`)
106+
const session = await connection.reserve()
107+
try {
108+
await session`BEGIN`
109+
for (const statement of source.split('--> statement-breakpoint'))
110+
if (statement.trim()) await session.unsafe(statement)
111+
await session`COMMIT`
112+
} finally {
113+
await session`ROLLBACK`
114+
await session`RESET statement_timeout`
115+
session.release()
116+
}
117+
}
118+
119+
async function ginState() {
120+
const [state] =
121+
await connection`SELECT c.oid, 'fastupdate=off' = ANY(c.reloptions) AS direct_writes,
122+
i.indisvalid, pending.pending_pages
123+
FROM pg_class c JOIN pg_index i ON i.indexrelid = c.oid
124+
CROSS JOIN LATERAL pgstatginindex(c.oid) pending
125+
WHERE c.oid = 'workspace_file_search_chunk_content_idx'::regclass`
126+
return state
127+
}
96128

97129
async function addFile(
98130
fileId: string,
@@ -110,10 +142,14 @@ describe('chunked workspace file search on PostgreSQL', () => {
110142
expect(build).not.toBeNull()
111143
const plan = planFileSearchIndex({ text, partial: false }, signal)
112144
const chunks = [...iterateFileSearchChunks(plan, signal)]
113-
for (let offset = 0; offset < chunks.length; offset += 100)
114-
expect(await appendFileSearchChunks(build!, chunks.slice(offset, offset + 100), signal)).toBe(
115-
true
116-
)
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)
117153
expect(
118154
await publishFileSearchBuild(
119155
build!,
@@ -140,6 +176,7 @@ describe('chunked workspace file search on PostgreSQL', () => {
140176
let captureQuery = false
141177

142178
beforeAll(async () => {
179+
await connection`CREATE EXTENSION IF NOT EXISTS pgstattuple`
143180
await connection`CREATE SCHEMA ${connection(schema)}`
144181
await connection`CREATE TABLE workspace (id text PRIMARY KEY)`
145182
await connection`CREATE TABLE workspace_files (id text PRIMARY KEY, workspace_id text REFERENCES workspace(id) ON DELETE CASCADE,
@@ -149,13 +186,9 @@ describe('chunked workspace file search on PostgreSQL', () => {
149186
'0313_puzzling_zodiak.sql',
150187
'0358_workspace_file_content_version_precision.sql',
151188
'0359_workspace_file_search_chunks.sql',
189+
ginWriteMigration,
152190
]) {
153-
const source = readFileSync(
154-
resolve(process.cwd(), '../../packages/db/migrations', migration),
155-
'utf8'
156-
).replaceAll('"public".', `"${schema}".`)
157-
for (const statement of source.split('--> statement-breakpoint'))
158-
if (statement.trim()) await connection.unsafe(statement)
191+
await applyMigration(migration)
159192
}
160193
database.current = drizzle(connection)
161194
database.search = drizzle(searchConnection, {
@@ -189,6 +222,88 @@ describe('chunked workspace file search on PostgreSQL', () => {
189222
}
190223
})
191224

225+
it('preserves search through disabling, draining, and replaying GIN pending-list maintenance', async () => {
226+
await connection`ALTER INDEX workspace_file_search_chunk_content_idx SET (fastupdate = on)`
227+
try {
228+
await index('heading\nold needle αβγ\ntail')
229+
const before = await ginState()
230+
expect(before.pending_pages).toBeGreaterThan(0)
231+
232+
/** Simulate an interrupted rollout after the storage option commits but before the drain. */
233+
await connection`ALTER INDEX workspace_file_search_chunk_content_idx SET (fastupdate = off)`
234+
await index('heading\nnew needle αβγ\ntail', await addFile('file-2'))
235+
expect((await ginState()).pending_pages).toBe(before.pending_pages)
236+
const expected = [
237+
{ fileId: 'file-1', lineNumber: 2 },
238+
{ fileId: 'file-2', lineNumber: 2 },
239+
]
240+
expect((await search('^(old|new) needle αβγ$', 'regex')).results).toMatchObject(expected)
241+
242+
for (let attempt = 0; attempt < 2; attempt++) {
243+
await applyMigration(ginWriteMigration)
244+
expect(await ginState()).toMatchObject({
245+
oid: before.oid,
246+
indisvalid: true,
247+
direct_writes: true,
248+
pending_pages: 0,
249+
})
250+
expect((await search('needle αβγ')).results).toMatchObject(expected)
251+
}
252+
await index('heading\nnew needle αβγ\ntail', await addFile('file-3'))
253+
expect((await ginState()).pending_pages).toBe(0)
254+
expect((await search('^(old|new) needle αβγ$', 'regex')).results).toHaveLength(3)
255+
} finally {
256+
await applyMigration(ginWriteMigration)
257+
}
258+
})
259+
260+
it('cancels a slow chunk statement without losing the connection or publishing partial content', async () => {
261+
const build = (await beginFileSearchBuild(revision))!
262+
const plan = planFileSearchIndex({ text: 'needle', partial: false }, signal)
263+
const chunks = [...iterateFileSearchChunks(plan, signal)]
264+
const writer = postgres(databaseUrl, {
265+
max: 1,
266+
prepare: false,
267+
connection: { search_path: `${schema},public` },
268+
})
269+
const original = database.current
270+
await connection`CREATE FUNCTION slow_chunk_insert() RETURNS trigger LANGUAGE plpgsql AS $$
271+
BEGIN PERFORM pg_sleep(11); RETURN NEW; END $$`
272+
await connection`CREATE TRIGGER slow_chunk_insert BEFORE INSERT ON workspace_file_search_chunk
273+
FOR EACH STATEMENT EXECUTE FUNCTION slow_chunk_insert()`
274+
try {
275+
database.current = drizzle(writer)
276+
const [before] = await writer`SELECT pg_backend_pid() AS pid`
277+
await expect(appendFileSearchChunks(build, chunks, signal)).rejects.toMatchObject({
278+
cause: { code: '57014' },
279+
})
280+
expect((await writer`SELECT pg_backend_pid() AS pid`)[0].pid).toBe(before.pid)
281+
expect(
282+
(await connection`SELECT count(*)::int AS count FROM workspace_file_search_chunk`)[0].count
283+
).toBe(0)
284+
expect((await search('needle')).results).toEqual([])
285+
} finally {
286+
database.current = original
287+
await writer.end()
288+
await connection`DROP TRIGGER slow_chunk_insert ON workspace_file_search_chunk`
289+
await connection`DROP FUNCTION slow_chunk_insert()`
290+
}
291+
expect(await appendFileSearchChunks(build, chunks, signal)).toBe(true)
292+
expect(
293+
await publishFileSearchBuild(
294+
build,
295+
{
296+
status: 'ready',
297+
chunkCount: chunks.length,
298+
lineCount: plan.lineCount,
299+
indexedBytes: plan.indexedBytes,
300+
},
301+
signal
302+
)
303+
).toBe(true)
304+
expect((await search('needle')).results).toMatchObject([{ fileId: 'file-1', lineNumber: 1 }])
305+
})
306+
192307
it('packs a million short lines without a million rows and bounds every stored value', async () => {
193308
await index('abc\n'.repeat(1_000_000))
194309
const [row] =

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

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,15 @@ export const FILE_SEARCH_CLEANUP_MIN_BATCH_MS =
4646
FILE_SEARCH_CLEANUP_BUDGET_MS / FILE_SEARCH_CLEANUP_MAX_BATCHES
4747
export const FILE_SEARCH_RECONCILE_INTERVAL_MS = 60 * 60 * 1000
4848
export const FILE_SEARCH_INSERT_BATCH_ROWS = 250
49-
export const FILE_SEARCH_INSERT_BATCH_BYTES = 1024 * 1024
49+
/** Direct GIN writes perform index work in each insert, so transactions use smaller byte batches. */
50+
export const FILE_SEARCH_INSERT_BATCH_BYTES = 128 * 1024
51+
52+
/** Index writes allow statement cancellation before the outer transaction terminates its session. */
53+
export const FILE_SEARCH_INDEX_TRANSACTION_LIMITS = {
54+
statementTimeout: 10 * 1000,
55+
lockTimeout: 5 * 1000,
56+
transactionTimeout: 30 * 1000,
57+
} as const
5058

5159
export const FILE_SEARCH_INDEX_GLOBAL_CONCURRENCY = 10
5260
export const FILE_SEARCH_INDEX_WORKSPACE_OUTSTANDING = 2

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

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import {
1515
FILE_SEARCH_CLEANUP_BUDGET_MS,
1616
FILE_SEARCH_CLEANUP_MAX_BATCHES,
1717
FILE_SEARCH_CLEANUP_MIN_BATCH_MS,
18+
FILE_SEARCH_INDEX_TRANSACTION_LIMITS,
1819
FILE_SEARCH_INSERT_BATCH_BYTES,
1920
FILE_SEARCH_INSERT_BATCH_ROWS,
2021
} from '@/lib/workspace-files/search/constants'
@@ -90,7 +91,7 @@ export async function beginFileSearchBuild(
9091
dispatchToken?: string
9192
): Promise<FileSearchBuild | null> {
9293
return db.transaction(async (tx) => {
93-
await configureFileSearchTransaction(tx)
94+
await configureFileSearchTransaction(tx, FILE_SEARCH_INDEX_TRANSACTION_LIMITS)
9495
if (!(await lockCurrentFile(tx, revision))) return null
9596
const [observed] = await tx
9697
.select({
@@ -154,7 +155,7 @@ export async function appendFileSearchChunks(
154155
throw new Error('File search insert batch exceeds its budget')
155156
}
156157
return db.transaction(async (tx) => {
157-
await configureFileSearchTransaction(tx)
158+
await configureFileSearchTransaction(tx, FILE_SEARCH_INDEX_TRANSACTION_LIMITS)
158159
if (!(await lockBuild(tx, build))) return false
159160
signal.throwIfAborted()
160161
await tx
@@ -177,7 +178,7 @@ export async function publishFileSearchBuild(
177178
signal: AbortSignal
178179
): Promise<boolean> {
179180
return db.transaction(async (tx) => {
180-
await configureFileSearchTransaction(tx)
181+
await configureFileSearchTransaction(tx, FILE_SEARCH_INDEX_TRANSACTION_LIMITS)
181182
signal.throwIfAborted()
182183
if (!(await lockCurrentFile(tx, build)) || !(await lockBuild(tx, build))) return false
183184
if (publication.status === 'ready') {
@@ -216,7 +217,7 @@ export async function failFileSearchRevision(
216217
dispatchToken?: string
217218
): Promise<void> {
218219
await db.transaction(async (tx) => {
219-
await configureFileSearchTransaction(tx)
220+
await configureFileSearchTransaction(tx, FILE_SEARCH_INDEX_TRANSACTION_LIMITS)
220221
if (!(await lockCurrentFile(tx, revision))) return
221222
const [state] = await tx
222223
.select()
Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,12 @@
1+
-- migration-safe: Metadata-only storage option; existing and new workers retain the same index and search semantics.
2+
ALTER INDEX "public"."workspace_file_search_chunk_content_idx" SET (fastupdate = off);
3+
--> statement-breakpoint
4+
-- Release the DDL lock before maintenance. Both operations are safe to replay after a partial migration.
5+
COMMIT;
6+
--> statement-breakpoint
7+
-- Disabling fastupdate does not flush existing pending entries. Drain them once without rebuilding the index.
8+
SET statement_timeout = '5min';
9+
--> statement-breakpoint
10+
SELECT pg_catalog.gin_clean_pending_list('"public"."workspace_file_search_chunk_content_idx"'::regclass);
11+
--> statement-breakpoint
12+
SET statement_timeout = 0;

0 commit comments

Comments
 (0)