From faa32e484e226ac5b4218766978ea21a8403dc8d Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 6 Oct 2026 18:18:03 -0700 Subject: [PATCH] improvement(file-search): only start the dispatcher when there is dispatch work --- .../workspace-file-search-dispatch/route.ts | 4 +- .../search/dispatcher.integration.ts | 113 +++++++++++++++++- .../lib/workspace-files/search/dispatcher.ts | 87 +++++++++++--- .../search/enqueue-dispatch.test.ts | 38 +++++- .../search/enqueue-dispatch.ts | 48 ++++++-- .../lib/workspace-files/search/index-state.ts | 8 +- 6 files changed, 255 insertions(+), 43 deletions(-) diff --git a/apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts b/apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts index 2ab25d6ba2f..f245a89906c 100644 --- a/apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts +++ b/apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts @@ -16,8 +16,8 @@ export const GET = withRouteHandler(async (request: NextRequest) => { try { const result = await enqueueWorkspaceFileSearchDispatch() - logger.info('Workspace file search dispatcher accepted', result) - return NextResponse.json({ success: true, triggered: true, ...result }, { status: 202 }) + if (result.triggered) logger.info('Workspace file search dispatcher accepted', result) + return NextResponse.json({ success: true, ...result }, { status: result.triggered ? 202 : 200 }) } catch (error) { logger.error('Workspace file search dispatcher enqueue failed', { error: toError(error).message, diff --git a/apps/sim/lib/workspace-files/search/dispatcher.integration.ts b/apps/sim/lib/workspace-files/search/dispatcher.integration.ts index e32b13aaccb..52b6c49cfad 100644 --- a/apps/sim/lib/workspace-files/search/dispatcher.integration.ts +++ b/apps/sim/lib/workspace-files/search/dispatcher.integration.ts @@ -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()), 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,9 +95,9 @@ 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) }) @@ -102,9 +105,11 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => { 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) => 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/) + }) + }) }) diff --git a/apps/sim/lib/workspace-files/search/dispatcher.ts b/apps/sim/lib/workspace-files/search/dispatcher.ts index 377380f0ed0..ddfbb9a991d 100644 --- a/apps/sim/lib/workspace-files/search/dispatcher.ts +++ b/apps/sim/lib/workspace-files/search/dispatcher.ts @@ -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 { .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 { + const [cursor] = await db + .select({ completedAt: workspaceFileSearchBackfill.completedAt }) + .from(workspaceFileSearchBackfill) + .where(eq(workspaceFileSearchBackfill.id, BACKFILL_CURSOR_ID)) + .limit(1) + 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 { diff --git a/apps/sim/lib/workspace-files/search/enqueue-dispatch.test.ts b/apps/sim/lib/workspace-files/search/enqueue-dispatch.test.ts index b89eb64f93d..a2b42a234c7 100644 --- a/apps/sim/lib/workspace-files/search/enqueue-dispatch.test.ts +++ b/apps/sim/lib/workspace-files/search/enqueue-dispatch.test.ts @@ -8,15 +8,14 @@ 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' @@ -24,15 +23,16 @@ import { enqueueWorkspaceFileSearchDispatch } from '@/lib/workspace-files/search 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) + }) }) diff --git a/apps/sim/lib/workspace-files/search/enqueue-dispatch.ts b/apps/sim/lib/workspace-files/search/enqueue-dispatch.ts index 750fd96d56a..d54e83b660b 100644 --- a/apps/sim/lib/workspace-files/search/enqueue-dispatch.ts +++ b/apps/sim/lib/workspace-files/search/enqueue-dispatch.ts @@ -1,14 +1,36 @@ +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' import { isTriggerDevEnabled } from '@/lib/core/config/env-flags' import { runDetached } from '@/lib/core/utils/background' import { FILE_SEARCH_DISPATCH_INTERVAL_MS, FILE_SEARCH_DISPATCH_MAX_DURATION_SECONDS, } from '@/lib/workspace-files/search/constants' -import { dispatchWorkspaceFileSearchIndexJobs } from '@/lib/workspace-files/search/dispatcher' +import { + dispatchWorkspaceFileSearchIndexJobs, + hasWorkspaceFileSearchDispatchWork, +} from '@/lib/workspace-files/search/dispatcher' +import type { workspaceFileSearchDispatchTask } from '@/background/workspace-file-search-dispatch' + +const logger = createLogger('WorkspaceFileSearchDispatchEnqueue') -export interface WorkspaceFileSearchDispatchEnqueueResult { - backend: 'trigger-dev' | 'inline' - jobId: string | null +export type WorkspaceFileSearchDispatchEnqueueResult = + | { triggered: true; backend: 'trigger-dev' | 'inline'; jobId: string | null } + | { triggered: false; backend: null; jobId: null } + +/** + * Starts no dispatcher run when there is nothing to dispatch. A failed check starts one anyway, so + * an unhealthy probe can delay indexing by at most the dispatcher's own failure, never strand it. + */ +async function hasDispatchWork(): Promise { + try { + return await hasWorkspaceFileSearchDispatchWork(new Date()) + } catch (error) { + logger.warn('Workspace file search dispatch work check failed; dispatching anyway', { + error: getErrorMessage(error), + }) + return true + } } /** @@ -16,17 +38,19 @@ export interface WorkspaceFileSearchDispatchEnqueueResult { * development-only and detaches from the HTTP response because the local server is long-lived. */ export async function enqueueWorkspaceFileSearchDispatch(): Promise { + if (!(await hasDispatchWork())) { + return { triggered: false, backend: null, jobId: null } + } + if (!isTriggerDevEnabled) { runDetached('workspace-file-search-dispatch', dispatchWorkspaceFileSearchIndexJobs) - return { backend: 'inline', jobId: null } + return { triggered: true, backend: 'inline', jobId: null } } - const [{ tasks }, { workspaceFileSearchDispatchTask }, { resolveTriggerRegion }] = - await Promise.all([ - import('@trigger.dev/sdk'), - import('@/background/workspace-file-search-dispatch'), - import('@/lib/core/async-jobs/region'), - ]) + const [{ tasks }, { resolveTriggerRegion }] = await Promise.all([ + import('@trigger.dev/sdk'), + import('@/lib/core/async-jobs/region'), + ]) const scheduleWindow = Math.floor(Date.now() / FILE_SEARCH_DISPATCH_INTERVAL_MS) const handle = await tasks.trigger( 'workspace-file-search-dispatch', @@ -39,5 +63,5 @@ export async function enqueueWorkspaceFileSearchDispatch(): Promise { const deadline = Date.now() + FILE_SEARCH_CLEANUP_BUDGET_MS @@ -277,7 +283,7 @@ export async function cleanupFileSearchBuilds(): Promise { const builds = await tx.execute<{ id: string }>(sql`SELECT id FROM workspace_file_search_build - WHERE expires_at <= now() ORDER BY expires_at, id LIMIT ${FILE_SEARCH_CLEANUP_BATCH_BUILDS} FOR UPDATE SKIP LOCKED`) + WHERE ${fileSearchBuildExpired} ORDER BY expires_at, id LIMIT ${FILE_SEARCH_CLEANUP_BATCH_BUILDS} FOR UPDATE SKIP LOCKED`) if (!builds.length) return null const buildIds = sql.join( builds.map((build) => sql`${build.id}`),