Skip to content

Commit b37ce85

Browse files
authored
fix(workspace-files): keep the file-search claim on the ordered pending index instead of a hash join (#8128)
1 parent 7c2213f commit b37ce85

2 files changed

Lines changed: 54 additions & 8 deletions

File tree

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

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -89,7 +89,8 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
8989
await connection`CREATE INDEX workspace_files_workspace_active_keyset_idx
9090
ON workspace_files (workspace_id, id)
9191
WHERE deleted_at IS NULL AND context = 'workspace' AND workspace_id IS NOT NULL`
92-
await connection`CREATE INDEX ON workspace_file_search_revision
92+
await connection`CREATE INDEX workspace_file_search_revision_pending_idx
93+
ON workspace_file_search_revision
9394
(workspace_id, updated_at, file_id, source_content_updated_at)
9495
WHERE status = 'pending' AND dispatched_at IS NULL`
9596
await connection`CREATE INDEX ON workspace_file_search_revision (workspace_id, dispatched_at)
@@ -263,6 +264,39 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
263264
expect(plan).toMatch(/Index Cond:.*ROW\(/)
264265
}, 30_000)
265266

267+
it('claims from the ordered pending index rather than a hash join over the backlog', async () => {
268+
await seedQueue('workspace-1', 10_000)
269+
await connection`ANALYZE workspace_files`
270+
await connection`ANALYZE workspace_file_search_revision`
271+
272+
statements.length = 0
273+
await prepareWorkspaceFileSearchDispatch()
274+
275+
const claim = statements.find((statement) =>
276+
statement.query.includes('FOR UPDATE OF search_index SKIP LOCKED')
277+
)
278+
expect(claim).toBeDefined()
279+
280+
const plan = await connection.begin(async (tx) => {
281+
/**
282+
* At fixture scale the planner already nests the join, so it is pinned to the choice
283+
* production makes when it misestimates the timestamp equi-join under `FOR UPDATE`. A join
284+
* spelling then hashes every file and every pending revision of the workspace and sorts the
285+
* whole backlog before the top-N cut; the correlated LATERAL cannot be flattened into that
286+
* join, so the claim stays an ordered walk of the pending index.
287+
*/
288+
await tx`SET LOCAL enable_nestloop = off`
289+
const rows = await tx.unsafe(`EXPLAIN ${claim?.query}`, claim?.params as never[])
290+
return rows.map((row: Record<string, unknown>) => row['QUERY PLAN']).join('\n')
291+
})
292+
293+
/** The locked candidate scan must be fed by the ordered index walk, not a sorted hash join. */
294+
expect(plan).toMatch(
295+
/LockRows[^\n]*\n\s*-> {2}Nested Loop[^\n]*\n\s*-> {2}Index Scan using workspace_file_search_revision_pending_idx/
296+
)
297+
expect(plan).not.toMatch(/Sort Key: search_index(_\d+)?\.updated_at/)
298+
}, 30_000)
299+
266300
it('fails on a locked backfill row and releases the dispatcher lock', async () => {
267301
let release = () => {}
268302
let locked = () => {}

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

Lines changed: 19 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -303,7 +303,16 @@ async function reapStaleClaims(tx: DbTransaction, now: Date): Promise<number> {
303303
return rows.length
304304
}
305305

306-
/** Probe each workspace's available slots and lock candidates before the update, skipping busy rows. */
306+
/**
307+
* Probe each workspace's available slots and lock candidates before the update, skipping busy rows.
308+
*
309+
* The live-file check is a correlated LATERAL with `LIMIT 1` rather than a join. `FOR UPDATE`
310+
* forbids parallel plans and the planner estimates the timestamp equi-join at about one row, so
311+
* a join becomes a hash join over every file and every pending revision of the workspace before
312+
* the top-N sort, which exceeds the statement timeout on a large backlog. A LATERAL with LIMIT
313+
* cannot be flattened into that join, so the claim stays an ordered walk of the pending index
314+
* that stops after the batch size.
315+
*/
307316
async function claimQueuedWorkspaceJobs(
308317
tx: DbTransaction,
309318
workspaceIds: readonly string[],
@@ -339,12 +348,15 @@ async function claimQueuedWorkspaceJobs(
339348
SELECT search_index.workspace_id, search_index.file_id,
340349
search_index.source_content_updated_at, search_index.updated_at
341350
FROM workspace_file_search_revision AS search_index
342-
INNER JOIN workspace_files AS file
343-
ON file.id = search_index.file_id
344-
AND file.workspace_id = search_index.workspace_id
345-
AND file.context = 'workspace'
346-
AND file.deleted_at IS NULL
347-
AND file.content_updated_at = search_index.source_content_updated_at
351+
CROSS JOIN LATERAL (
352+
SELECT 1 FROM workspace_files AS file
353+
WHERE file.id = search_index.file_id
354+
AND file.workspace_id = search_index.workspace_id
355+
AND file.context = 'workspace'
356+
AND file.deleted_at IS NULL
357+
AND file.content_updated_at = search_index.source_content_updated_at
358+
LIMIT 1
359+
) AS live_file
348360
WHERE search_index.workspace_id = selected.workspace_id
349361
AND search_index.status = 'pending' AND search_index.dispatched_at IS NULL
350362
ORDER BY search_index.updated_at, search_index.file_id, search_index.source_content_updated_at

0 commit comments

Comments
 (0)