Repository navigation
improvement(file-search): only start the dispatcher when there is dispatch work #8714
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -18,7 +18,8 @@ vi.mock('@/lib/workspace-files/search/indexing', () => ({ | |
| indexWorkspaceFileForSearch: vi.fn(), | ||
| markWorkspaceFileSearchIndexFailed: vi.fn(), | ||
| })) | ||
| vi.mock('@/lib/workspace-files/search/index-state', () => ({ | ||
| vi.mock('@/lib/workspace-files/search/index-state', async (importOriginal) => ({ | ||
| ...(await importOriginal<typeof import('@/lib/workspace-files/search/index-state')>()), | ||
| cleanupFileSearchBuilds: vi.fn().mockResolvedValue(0), | ||
| })) | ||
| vi.mock('@trigger.dev/sdk', () => ({ tasks: { batchTrigger: mocks.batchTrigger } })) | ||
|
|
@@ -30,9 +31,11 @@ import { | |
| FILE_SEARCH_BACKFILL_PAGE_SIZE, | ||
| FILE_SEARCH_DISPATCH_HANDOFF_MS, | ||
| FILE_SEARCH_INDEX_STALE_DISPATCH_MS, | ||
| FILE_SEARCH_RECONCILE_INTERVAL_MS, | ||
| } from '@/lib/workspace-files/search/constants' | ||
| import { | ||
| dispatchWorkspaceFileSearchIndexJobs, | ||
| hasWorkspaceFileSearchDispatchWork, | ||
| prepareWorkspaceFileSearchDispatch, | ||
| } from '@/lib/workspace-files/search/dispatcher' | ||
|
|
||
|
|
@@ -92,19 +95,21 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => { | |
| WHERE status = 'pending' AND dispatched_at IS NULL` | ||
| await connection`CREATE INDEX ON workspace_file_search_revision (workspace_id, dispatched_at) | ||
| WHERE status = 'pending' AND dispatched_at IS NOT NULL` | ||
| await connection`INSERT INTO workspace_file_search_backfill (id, updated_at) | ||
| VALUES ('workspace-file-search-chunks-v2', '2026-09-16 00:00:00')` | ||
| await connection`CREATE TABLE workspace_file_search_build (id text PRIMARY KEY, expires_at timestamp)` | ||
| await connection`CREATE INDEX workspace_file_search_build_cleanup_idx | ||
| ON workspace_file_search_build (expires_at, id) WHERE expires_at IS NOT NULL` | ||
| await connection`CREATE TABLE workspace_file_search_chunk (build_id text NOT NULL, ordinal integer NOT NULL, PRIMARY KEY(build_id, ordinal))` | ||
| database.current = drizzle(connection) | ||
| }) | ||
|
|
||
| beforeEach(async () => { | ||
| mocks.batchTrigger.mockReset() | ||
| await connection`DROP TRIGGER IF EXISTS slow_backfill ON workspace_file_search_backfill` | ||
| await connection`TRUNCATE workspace_files, workspace_file_search_revision, workspace_file_search_dispatch_queue` | ||
| await connection`UPDATE workspace_file_search_backfill | ||
| SET updated_at = '2026-09-16 00:00:00', completed_at = NULL, | ||
| await connection`TRUNCATE workspace_files, workspace_file_search_revision, | ||
| workspace_file_search_dispatch_queue, workspace_file_search_build` | ||
| await connection`INSERT INTO workspace_file_search_backfill (id, updated_at) | ||
| VALUES ('workspace-file-search-chunks-v2', '2026-09-16 00:00:00') | ||
| ON CONFLICT (id) DO UPDATE SET updated_at = EXCLUDED.updated_at, completed_at = NULL, | ||
| after_workspace_id = NULL, after_file_id = NULL` | ||
| }) | ||
|
|
||
|
|
@@ -544,4 +549,100 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => { | |
| { lock_timeout: '0', statement_timeout: '0', transaction_timeout: '0' }, | ||
| ]) | ||
| }) | ||
|
|
||
| describe('dispatch work check', () => { | ||
| const reconcileSeconds = FILE_SEARCH_RECONCILE_INTERVAL_MS / 1000 | ||
| const staleSeconds = FILE_SEARCH_INDEX_STALE_DISPATCH_MS / 1000 | ||
|
|
||
| /** | ||
| * Everything a deployment between file changes holds: a reconciled backfill, a published | ||
| * revision, a claim whose run is under way, a live build, and a build still inside its lease. | ||
| */ | ||
| async function seedIdleDeployment() { | ||
| await connection`UPDATE workspace_file_search_backfill | ||
| SET completed_at = now() - make_interval(secs => ${reconcileSeconds - 60})` | ||
| await connection`INSERT INTO workspace_files (id, workspace_id, context, content_updated_at) | ||
| VALUES ('ready-file', 'workspace-1', 'workspace', '2026-09-16'), | ||
| ('running-file', 'workspace-1', 'workspace', '2026-09-16')` | ||
| await connection`INSERT INTO workspace_file_search_revision | ||
| (file_id, workspace_id, source_content_updated_at, status, dispatched_at, handoff_expires_at) | ||
| VALUES ('ready-file', 'workspace-1', '2026-09-16', 'ready', NULL, NULL), | ||
| ('running-file', 'workspace-1', '2026-09-16', 'pending', | ||
| now() - make_interval(secs => ${staleSeconds - 60}), clock_timestamp() + interval '1 minute')` | ||
| await connection`INSERT INTO workspace_file_search_build (id, expires_at) | ||
| VALUES ('published-build', NULL), ('leased-build', now() + interval '1 minute')` | ||
| } | ||
|
|
||
| it('reports no work for an idle deployment, on which a dispatch pass does nothing', async () => { | ||
| await seedIdleDeployment() | ||
|
|
||
| await expect(hasWorkspaceFileSearchDispatchWork(new Date())).resolves.toBe(false) | ||
| await expect(prepareWorkspaceFileSearchDispatch()).resolves.toMatchObject({ | ||
| payloads: [], | ||
| backfilledFiles: 0, | ||
| reapedClaims: 0, | ||
| }) | ||
| }) | ||
|
|
||
| it.each([ | ||
| { | ||
| work: 'a backfill pass that never completed', | ||
| seed: () => connection`UPDATE workspace_file_search_backfill SET completed_at = NULL`, | ||
| }, | ||
| { | ||
| work: 'a missing backfill cursor', | ||
| seed: () => connection`DELETE FROM workspace_file_search_backfill`, | ||
| }, | ||
| { | ||
| work: 'a backfill reconcile that is due', | ||
| seed: () => connection`UPDATE workspace_file_search_backfill | ||
| SET completed_at = now() - make_interval(secs => ${reconcileSeconds})`, | ||
| }, | ||
| { | ||
| work: 'a queued workspace', | ||
| seed: () => connection`INSERT INTO workspace_file_search_dispatch_queue | ||
| (workspace_id, enqueued_at, updated_at) VALUES ('workspace-1', now(), now())`, | ||
| }, | ||
| { | ||
| work: 'an expired build', | ||
| seed: () => connection`UPDATE workspace_file_search_build SET expires_at = now() | ||
| WHERE id = 'leased-build'`, | ||
| }, | ||
| { | ||
| work: 'a claim past the stale-dispatch window', | ||
| seed: () => connection`UPDATE workspace_file_search_revision | ||
| SET dispatched_at = now() - make_interval(secs => ${staleSeconds + 60}) | ||
| WHERE file_id = 'running-file'`, | ||
| }, | ||
| { | ||
| work: 'a claim past its handoff deadline', | ||
| seed: () => connection`UPDATE workspace_file_search_revision | ||
| SET handoff_expires_at = clock_timestamp() - interval '1 millisecond' | ||
| WHERE file_id = 'running-file'`, | ||
| }, | ||
| ])('reports work for $work', async ({ seed }) => { | ||
| await seedIdleDeployment() | ||
| await seed() | ||
|
|
||
| await expect(hasWorkspaceFileSearchDispatchWork(new Date())).resolves.toBe(true) | ||
| }) | ||
|
|
||
| it('probes the cleanup and active-claim indexes rather than scanning', async () => { | ||
| await seedIdleDeployment() | ||
| statements.length = 0 | ||
| await expect(hasWorkspaceFileSearchDispatchWork(new Date())).resolves.toBe(false) | ||
|
|
||
| const probe = statements.find((statement) => statement.query.includes('"hasWork"')) | ||
| expect(probe).toBeDefined() | ||
| const plan = await connection.begin(async (tx) => { | ||
| await tx`SET LOCAL enable_seqscan = off` | ||
| const rows = await tx.unsafe(`EXPLAIN ${probe?.query}`, probe?.params as never[]) | ||
| return rows.map((row: Record<string, unknown>) => row['QUERY PLAN']).join('\n') | ||
| }) | ||
|
|
||
| expect(plan).toContain('workspace_file_search_build_cleanup_idx') | ||
| expect(plan).toMatch(/using workspace_file_search_revision_workspace_id_dispatched_at_idx/) | ||
| expect(plan).not.toMatch(/Seq Scan/) | ||
|
Comment on lines
+643
to
+645
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The new plan test expects a fixture index on |
||
| }) | ||
| }) | ||
| }) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -1,6 +1,7 @@ | ||
| import { db } from '@sim/db' | ||
| import { | ||
| workspaceFileSearchBackfill, | ||
| workspaceFileSearchBuild, | ||
| workspaceFileSearchDispatchQueue, | ||
| workspaceFileSearchRevision, | ||
| workspaceFiles, | ||
|
|
@@ -48,7 +49,10 @@ import { | |
| FILE_SEARCH_INDEX_WORKSPACE_OUTSTANDING, | ||
| FILE_SEARCH_RECONCILE_INTERVAL_MS, | ||
| } from '@/lib/workspace-files/search/constants' | ||
| import { cleanupFileSearchBuilds } from '@/lib/workspace-files/search/index-state' | ||
| import { | ||
| cleanupFileSearchBuilds, | ||
| fileSearchBuildExpired, | ||
| } from '@/lib/workspace-files/search/index-state' | ||
| import { | ||
| indexWorkspaceFileForSearch, | ||
| markWorkspaceFileSearchIndexFailed, | ||
|
|
@@ -166,6 +170,32 @@ async function enqueueWorkspaces( | |
| }) | ||
| } | ||
|
|
||
| /** | ||
| * Whether the backfill walk owes a page: the cursor has never completed a pass, or its last | ||
| * complete pass is at least {@link FILE_SEARCH_RECONCILE_INTERVAL_MS} old. | ||
| */ | ||
| function isBackfillPageDue(completedAt: Date | null, now: Date): boolean { | ||
| return !completedAt || now.getTime() - completedAt.getTime() >= FILE_SEARCH_RECONCILE_INTERVAL_MS | ||
| } | ||
|
|
||
| /** | ||
| * A claim {@link reapStaleClaims} releases: its handoff deadline has passed (in PostgreSQL time), or | ||
| * it has outlasted the stale-dispatch window. Served by `workspace_file_search_revision_active_idx`. | ||
| */ | ||
| function staleClaim(now: Date): SQL | undefined { | ||
| return and( | ||
| eq(workspaceFileSearchRevision.status, 'pending'), | ||
| isNotNull(workspaceFileSearchRevision.dispatchedAt), | ||
| or( | ||
| lt( | ||
| workspaceFileSearchRevision.dispatchedAt, | ||
| new Date(now.getTime() - FILE_SEARCH_INDEX_STALE_DISPATCH_MS) | ||
| ), | ||
| lte(workspaceFileSearchRevision.handoffExpiresAt, sql`clock_timestamp()`) | ||
| ) | ||
| ) | ||
| } | ||
|
|
||
| /** | ||
| * Seeds one page of the backfill that walks every live workspace file into the revision table. | ||
| * | ||
|
|
@@ -192,12 +222,7 @@ async function seedBackfillPage(tx: DbTransaction, now: Date): Promise<number> { | |
| .where(eq(workspaceFileSearchBackfill.id, BACKFILL_CURSOR_ID)) | ||
| .for('update') | ||
| .limit(1) | ||
| if ( | ||
| !cursor || | ||
| (cursor.completedAt && | ||
| now.getTime() - cursor.completedAt.getTime() < FILE_SEARCH_RECONCILE_INTERVAL_MS) | ||
| ) | ||
| return 0 | ||
| if (!cursor || !isBackfillPageDue(cursor.completedAt, now)) return 0 | ||
| const afterWorkspaceId = cursor.completedAt ? null : cursor.afterWorkspaceId | ||
| const afterFileId = cursor.completedAt ? null : cursor.afterFileId | ||
|
|
||
|
|
@@ -271,7 +296,6 @@ async function reapStaleClaims( | |
| tx: DbTransaction, | ||
| now: Date | ||
| ): Promise<{ reaped: number; abandoned: number }> { | ||
| const staleBefore = new Date(now.getTime() - FILE_SEARCH_INDEX_STALE_DISPATCH_MS) | ||
| const rows = await tx | ||
| .select({ | ||
| workspaceId: workspaceFileSearchRevision.workspaceId, | ||
|
|
@@ -291,16 +315,7 @@ async function reapStaleClaims( | |
| eq(workspaceFiles.contentUpdatedAt, workspaceFileSearchRevision.sourceContentUpdatedAt) | ||
| ) | ||
| ) | ||
| .where( | ||
| and( | ||
| eq(workspaceFileSearchRevision.status, 'pending'), | ||
| isNotNull(workspaceFileSearchRevision.dispatchedAt), | ||
| or( | ||
| lt(workspaceFileSearchRevision.dispatchedAt, staleBefore), | ||
| lte(workspaceFileSearchRevision.handoffExpiresAt, sql`clock_timestamp()`) | ||
| ) | ||
| ) | ||
| ) | ||
| .where(staleClaim(now)) | ||
| .orderBy(asc(workspaceFileSearchRevision.dispatchedAt), asc(workspaceFileSearchRevision.fileId)) | ||
| .limit(FILE_SEARCH_INDEX_STALE_REAP_LIMIT) | ||
| .for('update', { of: workspaceFileSearchRevision, skipLocked: true }) | ||
|
|
@@ -449,6 +464,42 @@ async function claimQueuedWorkspaceJobs( | |
| })) | ||
| } | ||
|
|
||
| /** | ||
| * Whether a dispatch pass at `now` has anything to do, so the per-minute cron can skip starting one | ||
| * on an idle deployment. Mirrors each phase of {@link dispatchWorkspaceFileSearchIndexJobs} with the | ||
| * same predicates: a backfill page is due, a build has expired for cleanup, a claim is stale for the | ||
| * reaper, or a workspace is queued for claiming. File writes, file deletes, and released claims | ||
| * all land in one of these through the `workspace_file_search_mark_pending` trigger or the | ||
| * dispatcher's own writes. | ||
| */ | ||
| export async function hasWorkspaceFileSearchDispatchWork(now: Date): Promise<boolean> { | ||
| const [cursor] = await db | ||
| .select({ completedAt: workspaceFileSearchBackfill.completedAt }) | ||
| .from(workspaceFileSearchBackfill) | ||
| .where(eq(workspaceFileSearchBackfill.id, BACKFILL_CURSOR_ID)) | ||
| .limit(1) | ||
|
Comment on lines
+475
to
+480
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The new database reads run before enqueueing without setting query or lock timeouts. If a table lock or slow query holds either read past the cron request's 60-second limit, |
||
| if (!cursor || isBackfillPageDue(cursor.completedAt, now)) return true | ||
|
|
||
| const queuedWorkspace = db | ||
| .select({ workspaceId: workspaceFileSearchDispatchQueue.workspaceId }) | ||
| .from(workspaceFileSearchDispatchQueue) | ||
| .limit(1) | ||
| const expiredBuild = db | ||
| .select({ id: workspaceFileSearchBuild.id }) | ||
| .from(workspaceFileSearchBuild) | ||
| .where(fileSearchBuildExpired) | ||
| .limit(1) | ||
| const staleClaimRow = db | ||
| .select({ fileId: workspaceFileSearchRevision.fileId }) | ||
| .from(workspaceFileSearchRevision) | ||
| .where(staleClaim(now)) | ||
| .limit(1) | ||
| const [probe] = await db.execute<{ hasWork: boolean }>( | ||
| sql`SELECT ${or(exists(queuedWorkspace), exists(expiredBuild), exists(staleClaimRow))} AS "hasWork"` | ||
| ) | ||
| return probe?.hasWork === true | ||
| } | ||
|
|
||
| export async function prepareWorkspaceFileSearchDispatch( | ||
| maxOutstanding = FILE_SEARCH_INDEX_MAX_OUTSTANDING | ||
| ): Promise<PreparedDispatch> { | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -8,31 +8,31 @@ import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vites | |
|
|
||
| const mocks = vi.hoisted(() => ({ | ||
| dispatch: vi.fn(), | ||
| hasWork: vi.fn(), | ||
| })) | ||
|
|
||
| vi.mock('@/background/workspace-file-search-dispatch', () => ({ | ||
| workspaceFileSearchDispatchTask: {}, | ||
| })) | ||
| vi.mock('@/lib/core/async-jobs/region', () => asyncJobsRegionMock) | ||
| vi.mock('@/lib/core/utils/background', () => backgroundTaskMock) | ||
| vi.mock('@/lib/workspace-files/search/dispatcher', () => ({ | ||
| dispatchWorkspaceFileSearchIndexJobs: mocks.dispatch, | ||
| hasWorkspaceFileSearchDispatchWork: mocks.hasWork, | ||
| })) | ||
|
|
||
| import { tasks } from '@trigger.dev/sdk' | ||
| import { enqueueWorkspaceFileSearchDispatch } from '@/lib/workspace-files/search/enqueue-dispatch' | ||
|
|
||
| const mockTrigger = vi.mocked(tasks.trigger) | ||
|
|
||
| setEnvFlags({ isTriggerDevEnabled: true }) | ||
| afterAll(resetEnvFlagsMock) | ||
|
|
||
| describe('workspace file search dispatcher enqueue', () => { | ||
| beforeEach(() => { | ||
| vi.useFakeTimers() | ||
| vi.setSystemTime(new Date('2026-08-29T12:34:45.000Z')) | ||
| setEnvFlags({ isTriggerDevEnabled: true }) | ||
| asyncJobsRegionMockFns.mockResolveTriggerRegion.mockResolvedValue('us-east-1') | ||
| mockTrigger.mockResolvedValue({ id: 'run-1' }) | ||
| mocks.hasWork.mockResolvedValue(true) | ||
| }) | ||
|
|
||
| afterEach(() => { | ||
|
|
@@ -41,6 +41,7 @@ describe('workspace file search dispatcher enqueue', () => { | |
|
|
||
| it('waits only for durable Trigger.dev acceptance and does not run the dispatcher inline', async () => { | ||
| await expect(enqueueWorkspaceFileSearchDispatch()).resolves.toEqual({ | ||
| triggered: true, | ||
| backend: 'trigger-dev', | ||
| jobId: 'run-1', | ||
| }) | ||
|
|
@@ -55,4 +56,33 @@ describe('workspace file search dispatcher enqueue', () => { | |
| expect(mocks.dispatch).not.toHaveBeenCalled() | ||
| expect(backgroundTaskMockFns.mockRunDetached).not.toHaveBeenCalled() | ||
| }) | ||
|
|
||
| it.each([ | ||
| { backend: 'trigger-dev', triggerDevEnabled: true }, | ||
| { backend: 'inline', triggerDevEnabled: false }, | ||
| ])( | ||
| 'starts no $backend dispatcher run when there is no dispatch work', | ||
| async ({ triggerDevEnabled }) => { | ||
| setEnvFlags({ isTriggerDevEnabled: triggerDevEnabled }) | ||
| mocks.hasWork.mockResolvedValue(false) | ||
|
|
||
| await expect(enqueueWorkspaceFileSearchDispatch()).resolves.toEqual({ | ||
| triggered: false, | ||
| backend: null, | ||
| jobId: null, | ||
| }) | ||
| expect(mockTrigger).not.toHaveBeenCalled() | ||
| expect(backgroundTaskMockFns.mockRunDetached).not.toHaveBeenCalled() | ||
| } | ||
| ) | ||
|
|
||
| it('still starts a dispatcher run when the work check fails', async () => { | ||
| mocks.hasWork.mockRejectedValue(new Error('statement timeout')) | ||
|
|
||
| await expect(enqueueWorkspaceFileSearchDispatch()).resolves.toMatchObject({ | ||
| triggered: true, | ||
| backend: 'trigger-dev', | ||
| }) | ||
| expect(mockTrigger).toHaveBeenCalledTimes(1) | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The new failure-path test proves dispatch by asserting that Context Used: CLAUDE.md (source) Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time! |
||
| }) | ||
| }) | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
The new partial mock uses
importOriginal. The testing guide explicitly says never to usevi.importActual()orimportOriginalto build a partial mock; use the central mock instead. This loads the real module and its imports just to retainfileSearchBuildExpired. Keep the real module unmocked for this integration test, or use a faithful shared mock. This repository requirement must be satisfied before merging.Context Used: CLAUDE.md (source)
Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!