Skip to content

Commit 9950763

Browse files
committed
fix(file-search): seek the backfill cursor instead of rescanning each page
The hourly backfill walks every live workspace file by `(workspace_id, id)`, but no index supplied that order under its predicate, so each page sorted the whole remaining set and the dispatcher's 10s statement timeout aborted the transaction before any page committed. Adds the matching partial index and compares the cursor row-wise. The previous `workspace_id > :ws OR (workspace_id = :ws AND id > :id)` spelling is only ever an index filter, never an index condition, so even with the index each page restarted at the low end and rescanned every page before it. On a prod-shaped fixture (5.5M files, 94k live) a full 95-page walk goes from 786ms to 42ms; the index alone accounts for 786ms -> 267ms and the row-wise cursor for the rest.
1 parent 431966e commit 9950763

6 files changed

Lines changed: 27846 additions & 12 deletions

File tree

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

Lines changed: 42 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ vi.mock('@trigger.dev/sdk', () => ({ tasks: { batchTrigger: mocks.batchTrigger }
2525
vi.mock('@/lib/core/config/env-flags', () => ({ isTriggerDevEnabled: true }))
2626
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
2727

28+
import { FILE_SEARCH_BACKFILL_PAGE_SIZE } from '@/lib/workspace-files/search/constants'
2829
import {
2930
dispatchWorkspaceFileSearchIndexJobs,
3031
prepareWorkspaceFileSearchDispatch,
@@ -64,14 +65,19 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
6465
deleted_at timestamp, content_updated_at timestamp NOT NULL
6566
)`
6667
await connection`CREATE TABLE workspace_file_search_revision (
67-
file_id text NOT NULL, workspace_id text NOT NULL, source_content_updated_at timestamp NOT NULL,
68-
status text NOT NULL, dispatched_at timestamp, updated_at timestamp NOT NULL,
69-
PRIMARY KEY (file_id, source_content_updated_at)
68+
file_id text PRIMARY KEY, workspace_id text NOT NULL,
69+
source_content_updated_at timestamp NOT NULL, status text NOT NULL DEFAULT 'pending',
70+
build_id text, failure_reason text, line_count integer NOT NULL DEFAULT 0,
71+
indexed_bytes integer NOT NULL DEFAULT 0, chunk_count integer NOT NULL DEFAULT 0,
72+
dispatched_at timestamp, updated_at timestamp NOT NULL DEFAULT now()
7073
)`
7174
await connection`CREATE TABLE workspace_file_search_dispatch_queue (
7275
workspace_id text PRIMARY KEY, enqueued_at timestamp NOT NULL,
7376
updated_at timestamp NOT NULL, last_dispatched_at timestamp
7477
)`
78+
await connection`CREATE INDEX workspace_files_workspace_active_keyset_idx
79+
ON workspace_files (workspace_id, id)
80+
WHERE deleted_at IS NULL AND context = 'workspace' AND workspace_id IS NOT NULL`
7581
await connection`CREATE INDEX ON workspace_file_search_revision
7682
(workspace_id, updated_at, file_id, source_content_updated_at)
7783
WHERE status = 'pending' AND dispatched_at IS NULL`
@@ -89,7 +95,8 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
8995
await connection`DROP TRIGGER IF EXISTS slow_backfill ON workspace_file_search_backfill`
9096
await connection`TRUNCATE workspace_files, workspace_file_search_revision, workspace_file_search_dispatch_queue`
9197
await connection`UPDATE workspace_file_search_backfill
92-
SET updated_at = '2026-09-16 00:00:00', completed_at = NULL`
98+
SET updated_at = '2026-09-16 00:00:00', completed_at = NULL,
99+
after_workspace_id = NULL, after_file_id = NULL`
93100
})
94101

95102
afterAll(async () => {
@@ -176,6 +183,37 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
176183
expect(row.active).toBe(100)
177184
})
178185

186+
it('walks every live workspace file exactly once across backfill pages', async () => {
187+
const files = 2 * FILE_SEARCH_BACKFILL_PAGE_SIZE + FILE_SEARCH_BACKFILL_PAGE_SIZE / 2
188+
await connection`INSERT INTO workspace_files (id, workspace_id, context, content_updated_at)
189+
SELECT md5(n::text), 'workspace-' || lpad((n % 7)::text, 2, '0'), 'workspace', '2026-09-16'
190+
FROM generate_series(1, ${files}) n`
191+
await connection`INSERT INTO workspace_files
192+
(id, workspace_id, context, deleted_at, content_updated_at)
193+
VALUES ('skipped-deleted', 'workspace-00', 'workspace', now(), '2026-09-16')`
194+
await connection`INSERT INTO workspace_files (id, workspace_id, context, content_updated_at)
195+
VALUES ('skipped-context', 'workspace-00', 'execution', '2026-09-16')`
196+
197+
let pages = 0
198+
for (;;) {
199+
await prepareWorkspaceFileSearchDispatch()
200+
pages += 1
201+
const [cursor] = await connection`SELECT completed_at FROM workspace_file_search_backfill`
202+
if (cursor.completed_at) break
203+
expect(pages).toBeLessThanOrEqual(files)
204+
}
205+
206+
expect(pages).toBe(Math.ceil(files / FILE_SEARCH_BACKFILL_PAGE_SIZE))
207+
const [seeded] = await connection`SELECT count(*)::int AS total,
208+
count(DISTINCT file_id)::int AS distinct_files FROM workspace_file_search_revision`
209+
expect(seeded).toEqual({ total: files, distinct_files: files })
210+
const [skipped] = await connection`SELECT count(*)::int AS missed FROM workspace_files file
211+
WHERE file.context = 'workspace' AND file.deleted_at IS NULL
212+
AND NOT EXISTS (SELECT 1 FROM workspace_file_search_revision revision
213+
WHERE revision.file_id = file.id)`
214+
expect(skipped.missed).toBe(0)
215+
}, 30_000)
216+
179217
it('fails on a locked backfill row and releases the dispatcher lock', async () => {
180218
let release = () => {}
181219
let locked = () => {}

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

Lines changed: 10 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,6 @@ import {
1313
asc,
1414
eq,
1515
exists,
16-
gt,
1716
inArray,
1817
isNotNull,
1918
isNull,
@@ -154,6 +153,15 @@ async function enqueueWorkspaces(
154153
})
155154
}
156155

156+
/**
157+
* Seeds one page of the backfill that walks every live workspace file into the revision table.
158+
*
159+
* The cursor is compared row-wise so the walk seeks straight to it using
160+
* `workspace_files_workspace_active_keyset_idx`. The equivalent
161+
* `workspace_id > :ws OR (workspace_id = :ws AND id > :id)` spelling is not something the planner
162+
* can turn into an index condition: it stays a filter, so every page rescans the pages before it
163+
* and the walk degrades to O(files^2) until it exceeds the dispatch statement timeout.
164+
*/
157165
async function seedBackfillPage(tx: DbTransaction, now: Date): Promise<number> {
158166
await tx
159167
.insert(workspaceFileSearchBackfill)
@@ -188,13 +196,7 @@ async function seedBackfillPage(tx: DbTransaction, now: Date): Promise<number> {
188196
isNull(workspaceFiles.deletedAt),
189197
isNotNull(workspaceFiles.workspaceId),
190198
afterWorkspaceId && afterFileId
191-
? or(
192-
gt(workspaceFiles.workspaceId, afterWorkspaceId),
193-
and(
194-
eq(workspaceFiles.workspaceId, afterWorkspaceId),
195-
gt(workspaceFiles.id, afterFileId)
196-
)
197-
)
199+
? sql`(${workspaceFiles.workspaceId}, ${workspaceFiles.id}) > (${afterWorkspaceId}, ${afterFileId})`
198200
: undefined
199201
)
200202
)
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
COMMIT;--> statement-breakpoint
2+
SET lock_timeout = 0;--> statement-breakpoint
3+
-- migration-safe: replay replaces only this new index to recover an interrupted concurrent build; existing indexes remain available.
4+
DROP INDEX CONCURRENTLY IF EXISTS "workspace_files_workspace_active_keyset_idx";--> statement-breakpoint
5+
CREATE INDEX CONCURRENTLY IF NOT EXISTS "workspace_files_workspace_active_keyset_idx" ON "workspace_files" USING btree ("workspace_id","id") WHERE "workspace_files"."deleted_at" IS NULL AND "workspace_files"."context" = 'workspace' AND "workspace_files"."workspace_id" IS NOT NULL;--> statement-breakpoint
6+
SET lock_timeout = '5s';

0 commit comments

Comments
 (0)