Skip to content

Commit faa32e4

Browse files
committed
improvement(file-search): only start the dispatcher when there is dispatch work
1 parent 8fc4c1a commit faa32e4

6 files changed

Lines changed: 255 additions & 43 deletions

File tree

‎apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -16,8 +16,8 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
1616

1717
try {
1818
const result = await enqueueWorkspaceFileSearchDispatch()
19-
logger.info('Workspace file search dispatcher accepted', result)
20-
return NextResponse.json({ success: true, triggered: true, ...result }, { status: 202 })
19+
if (result.triggered) logger.info('Workspace file search dispatcher accepted', result)
20+
return NextResponse.json({ success: true, ...result }, { status: result.triggered ? 202 : 200 })
2121
} catch (error) {
2222
logger.error('Workspace file search dispatcher enqueue failed', {
2323
error: toError(error).message,

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

Lines changed: 107 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,8 @@ vi.mock('@/lib/workspace-files/search/indexing', () => ({
1818
indexWorkspaceFileForSearch: vi.fn(),
1919
markWorkspaceFileSearchIndexFailed: vi.fn(),
2020
}))
21-
vi.mock('@/lib/workspace-files/search/index-state', () => ({
21+
vi.mock('@/lib/workspace-files/search/index-state', async (importOriginal) => ({
22+
...(await importOriginal<typeof import('@/lib/workspace-files/search/index-state')>()),
2223
cleanupFileSearchBuilds: vi.fn().mockResolvedValue(0),
2324
}))
2425
vi.mock('@trigger.dev/sdk', () => ({ tasks: { batchTrigger: mocks.batchTrigger } }))
@@ -30,9 +31,11 @@ import {
3031
FILE_SEARCH_BACKFILL_PAGE_SIZE,
3132
FILE_SEARCH_DISPATCH_HANDOFF_MS,
3233
FILE_SEARCH_INDEX_STALE_DISPATCH_MS,
34+
FILE_SEARCH_RECONCILE_INTERVAL_MS,
3335
} from '@/lib/workspace-files/search/constants'
3436
import {
3537
dispatchWorkspaceFileSearchIndexJobs,
38+
hasWorkspaceFileSearchDispatchWork,
3639
prepareWorkspaceFileSearchDispatch,
3740
} from '@/lib/workspace-files/search/dispatcher'
3841

@@ -92,19 +95,21 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
9295
WHERE status = 'pending' AND dispatched_at IS NULL`
9396
await connection`CREATE INDEX ON workspace_file_search_revision (workspace_id, dispatched_at)
9497
WHERE status = 'pending' AND dispatched_at IS NOT NULL`
95-
await connection`INSERT INTO workspace_file_search_backfill (id, updated_at)
96-
VALUES ('workspace-file-search-chunks-v2', '2026-09-16 00:00:00')`
9798
await connection`CREATE TABLE workspace_file_search_build (id text PRIMARY KEY, expires_at timestamp)`
99+
await connection`CREATE INDEX workspace_file_search_build_cleanup_idx
100+
ON workspace_file_search_build (expires_at, id) WHERE expires_at IS NOT NULL`
98101
await connection`CREATE TABLE workspace_file_search_chunk (build_id text NOT NULL, ordinal integer NOT NULL, PRIMARY KEY(build_id, ordinal))`
99102
database.current = drizzle(connection)
100103
})
101104

102105
beforeEach(async () => {
103106
mocks.batchTrigger.mockReset()
104107
await connection`DROP TRIGGER IF EXISTS slow_backfill ON workspace_file_search_backfill`
105-
await connection`TRUNCATE workspace_files, workspace_file_search_revision, workspace_file_search_dispatch_queue`
106-
await connection`UPDATE workspace_file_search_backfill
107-
SET updated_at = '2026-09-16 00:00:00', completed_at = NULL,
108+
await connection`TRUNCATE workspace_files, workspace_file_search_revision,
109+
workspace_file_search_dispatch_queue, workspace_file_search_build`
110+
await connection`INSERT INTO workspace_file_search_backfill (id, updated_at)
111+
VALUES ('workspace-file-search-chunks-v2', '2026-09-16 00:00:00')
112+
ON CONFLICT (id) DO UPDATE SET updated_at = EXCLUDED.updated_at, completed_at = NULL,
108113
after_workspace_id = NULL, after_file_id = NULL`
109114
})
110115

@@ -544,4 +549,100 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
544549
{ lock_timeout: '0', statement_timeout: '0', transaction_timeout: '0' },
545550
])
546551
})
552+
553+
describe('dispatch work check', () => {
554+
const reconcileSeconds = FILE_SEARCH_RECONCILE_INTERVAL_MS / 1000
555+
const staleSeconds = FILE_SEARCH_INDEX_STALE_DISPATCH_MS / 1000
556+
557+
/**
558+
* Everything a deployment between file changes holds: a reconciled backfill, a published
559+
* revision, a claim whose run is under way, a live build, and a build still inside its lease.
560+
*/
561+
async function seedIdleDeployment() {
562+
await connection`UPDATE workspace_file_search_backfill
563+
SET completed_at = now() - make_interval(secs => ${reconcileSeconds - 60})`
564+
await connection`INSERT INTO workspace_files (id, workspace_id, context, content_updated_at)
565+
VALUES ('ready-file', 'workspace-1', 'workspace', '2026-09-16'),
566+
('running-file', 'workspace-1', 'workspace', '2026-09-16')`
567+
await connection`INSERT INTO workspace_file_search_revision
568+
(file_id, workspace_id, source_content_updated_at, status, dispatched_at, handoff_expires_at)
569+
VALUES ('ready-file', 'workspace-1', '2026-09-16', 'ready', NULL, NULL),
570+
('running-file', 'workspace-1', '2026-09-16', 'pending',
571+
now() - make_interval(secs => ${staleSeconds - 60}), clock_timestamp() + interval '1 minute')`
572+
await connection`INSERT INTO workspace_file_search_build (id, expires_at)
573+
VALUES ('published-build', NULL), ('leased-build', now() + interval '1 minute')`
574+
}
575+
576+
it('reports no work for an idle deployment, on which a dispatch pass does nothing', async () => {
577+
await seedIdleDeployment()
578+
579+
await expect(hasWorkspaceFileSearchDispatchWork(new Date())).resolves.toBe(false)
580+
await expect(prepareWorkspaceFileSearchDispatch()).resolves.toMatchObject({
581+
payloads: [],
582+
backfilledFiles: 0,
583+
reapedClaims: 0,
584+
})
585+
})
586+
587+
it.each([
588+
{
589+
work: 'a backfill pass that never completed',
590+
seed: () => connection`UPDATE workspace_file_search_backfill SET completed_at = NULL`,
591+
},
592+
{
593+
work: 'a missing backfill cursor',
594+
seed: () => connection`DELETE FROM workspace_file_search_backfill`,
595+
},
596+
{
597+
work: 'a backfill reconcile that is due',
598+
seed: () => connection`UPDATE workspace_file_search_backfill
599+
SET completed_at = now() - make_interval(secs => ${reconcileSeconds})`,
600+
},
601+
{
602+
work: 'a queued workspace',
603+
seed: () => connection`INSERT INTO workspace_file_search_dispatch_queue
604+
(workspace_id, enqueued_at, updated_at) VALUES ('workspace-1', now(), now())`,
605+
},
606+
{
607+
work: 'an expired build',
608+
seed: () => connection`UPDATE workspace_file_search_build SET expires_at = now()
609+
WHERE id = 'leased-build'`,
610+
},
611+
{
612+
work: 'a claim past the stale-dispatch window',
613+
seed: () => connection`UPDATE workspace_file_search_revision
614+
SET dispatched_at = now() - make_interval(secs => ${staleSeconds + 60})
615+
WHERE file_id = 'running-file'`,
616+
},
617+
{
618+
work: 'a claim past its handoff deadline',
619+
seed: () => connection`UPDATE workspace_file_search_revision
620+
SET handoff_expires_at = clock_timestamp() - interval '1 millisecond'
621+
WHERE file_id = 'running-file'`,
622+
},
623+
])('reports work for $work', async ({ seed }) => {
624+
await seedIdleDeployment()
625+
await seed()
626+
627+
await expect(hasWorkspaceFileSearchDispatchWork(new Date())).resolves.toBe(true)
628+
})
629+
630+
it('probes the cleanup and active-claim indexes rather than scanning', async () => {
631+
await seedIdleDeployment()
632+
statements.length = 0
633+
await expect(hasWorkspaceFileSearchDispatchWork(new Date())).resolves.toBe(false)
634+
635+
const probe = statements.find((statement) => statement.query.includes('"hasWork"'))
636+
expect(probe).toBeDefined()
637+
const plan = await connection.begin(async (tx) => {
638+
await tx`SET LOCAL enable_seqscan = off`
639+
const rows = await tx.unsafe(`EXPLAIN ${probe?.query}`, probe?.params as never[])
640+
return rows.map((row: Record<string, unknown>) => row['QUERY PLAN']).join('\n')
641+
})
642+
643+
expect(plan).toContain('workspace_file_search_build_cleanup_idx')
644+
expect(plan).toMatch(/using workspace_file_search_revision_workspace_id_dispatched_at_idx/)
645+
expect(plan).not.toMatch(/Seq Scan/)
646+
})
647+
})
547648
})

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

Lines changed: 69 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { db } from '@sim/db'
22
import {
33
workspaceFileSearchBackfill,
4+
workspaceFileSearchBuild,
45
workspaceFileSearchDispatchQueue,
56
workspaceFileSearchRevision,
67
workspaceFiles,
@@ -48,7 +49,10 @@ import {
4849
FILE_SEARCH_INDEX_WORKSPACE_OUTSTANDING,
4950
FILE_SEARCH_RECONCILE_INTERVAL_MS,
5051
} from '@/lib/workspace-files/search/constants'
51-
import { cleanupFileSearchBuilds } from '@/lib/workspace-files/search/index-state'
52+
import {
53+
cleanupFileSearchBuilds,
54+
fileSearchBuildExpired,
55+
} from '@/lib/workspace-files/search/index-state'
5256
import {
5357
indexWorkspaceFileForSearch,
5458
markWorkspaceFileSearchIndexFailed,
@@ -166,6 +170,32 @@ async function enqueueWorkspaces(
166170
})
167171
}
168172

173+
/**
174+
* Whether the backfill walk owes a page: the cursor has never completed a pass, or its last
175+
* complete pass is at least {@link FILE_SEARCH_RECONCILE_INTERVAL_MS} old.
176+
*/
177+
function isBackfillPageDue(completedAt: Date | null, now: Date): boolean {
178+
return !completedAt || now.getTime() - completedAt.getTime() >= FILE_SEARCH_RECONCILE_INTERVAL_MS
179+
}
180+
181+
/**
182+
* A claim {@link reapStaleClaims} releases: its handoff deadline has passed (in PostgreSQL time), or
183+
* it has outlasted the stale-dispatch window. Served by `workspace_file_search_revision_active_idx`.
184+
*/
185+
function staleClaim(now: Date): SQL | undefined {
186+
return and(
187+
eq(workspaceFileSearchRevision.status, 'pending'),
188+
isNotNull(workspaceFileSearchRevision.dispatchedAt),
189+
or(
190+
lt(
191+
workspaceFileSearchRevision.dispatchedAt,
192+
new Date(now.getTime() - FILE_SEARCH_INDEX_STALE_DISPATCH_MS)
193+
),
194+
lte(workspaceFileSearchRevision.handoffExpiresAt, sql`clock_timestamp()`)
195+
)
196+
)
197+
}
198+
169199
/**
170200
* Seeds one page of the backfill that walks every live workspace file into the revision table.
171201
*
@@ -192,12 +222,7 @@ async function seedBackfillPage(tx: DbTransaction, now: Date): Promise<number> {
192222
.where(eq(workspaceFileSearchBackfill.id, BACKFILL_CURSOR_ID))
193223
.for('update')
194224
.limit(1)
195-
if (
196-
!cursor ||
197-
(cursor.completedAt &&
198-
now.getTime() - cursor.completedAt.getTime() < FILE_SEARCH_RECONCILE_INTERVAL_MS)
199-
)
200-
return 0
225+
if (!cursor || !isBackfillPageDue(cursor.completedAt, now)) return 0
201226
const afterWorkspaceId = cursor.completedAt ? null : cursor.afterWorkspaceId
202227
const afterFileId = cursor.completedAt ? null : cursor.afterFileId
203228

@@ -271,7 +296,6 @@ async function reapStaleClaims(
271296
tx: DbTransaction,
272297
now: Date
273298
): Promise<{ reaped: number; abandoned: number }> {
274-
const staleBefore = new Date(now.getTime() - FILE_SEARCH_INDEX_STALE_DISPATCH_MS)
275299
const rows = await tx
276300
.select({
277301
workspaceId: workspaceFileSearchRevision.workspaceId,
@@ -291,16 +315,7 @@ async function reapStaleClaims(
291315
eq(workspaceFiles.contentUpdatedAt, workspaceFileSearchRevision.sourceContentUpdatedAt)
292316
)
293317
)
294-
.where(
295-
and(
296-
eq(workspaceFileSearchRevision.status, 'pending'),
297-
isNotNull(workspaceFileSearchRevision.dispatchedAt),
298-
or(
299-
lt(workspaceFileSearchRevision.dispatchedAt, staleBefore),
300-
lte(workspaceFileSearchRevision.handoffExpiresAt, sql`clock_timestamp()`)
301-
)
302-
)
303-
)
318+
.where(staleClaim(now))
304319
.orderBy(asc(workspaceFileSearchRevision.dispatchedAt), asc(workspaceFileSearchRevision.fileId))
305320
.limit(FILE_SEARCH_INDEX_STALE_REAP_LIMIT)
306321
.for('update', { of: workspaceFileSearchRevision, skipLocked: true })
@@ -449,6 +464,42 @@ async function claimQueuedWorkspaceJobs(
449464
}))
450465
}
451466

467+
/**
468+
* Whether a dispatch pass at `now` has anything to do, so the per-minute cron can skip starting one
469+
* on an idle deployment. Mirrors each phase of {@link dispatchWorkspaceFileSearchIndexJobs} with the
470+
* same predicates: a backfill page is due, a build has expired for cleanup, a claim is stale for the
471+
* reaper, or a workspace is queued for claiming. File writes, file deletes, and released claims
472+
* all land in one of these through the `workspace_file_search_mark_pending` trigger or the
473+
* dispatcher's own writes.
474+
*/
475+
export async function hasWorkspaceFileSearchDispatchWork(now: Date): Promise<boolean> {
476+
const [cursor] = await db
477+
.select({ completedAt: workspaceFileSearchBackfill.completedAt })
478+
.from(workspaceFileSearchBackfill)
479+
.where(eq(workspaceFileSearchBackfill.id, BACKFILL_CURSOR_ID))
480+
.limit(1)
481+
if (!cursor || isBackfillPageDue(cursor.completedAt, now)) return true
482+
483+
const queuedWorkspace = db
484+
.select({ workspaceId: workspaceFileSearchDispatchQueue.workspaceId })
485+
.from(workspaceFileSearchDispatchQueue)
486+
.limit(1)
487+
const expiredBuild = db
488+
.select({ id: workspaceFileSearchBuild.id })
489+
.from(workspaceFileSearchBuild)
490+
.where(fileSearchBuildExpired)
491+
.limit(1)
492+
const staleClaimRow = db
493+
.select({ fileId: workspaceFileSearchRevision.fileId })
494+
.from(workspaceFileSearchRevision)
495+
.where(staleClaim(now))
496+
.limit(1)
497+
const [probe] = await db.execute<{ hasWork: boolean }>(
498+
sql`SELECT ${or(exists(queuedWorkspace), exists(expiredBuild), exists(staleClaimRow))} AS "hasWork"`
499+
)
500+
return probe?.hasWork === true
501+
}
502+
452503
export async function prepareWorkspaceFileSearchDispatch(
453504
maxOutstanding = FILE_SEARCH_INDEX_MAX_OUTSTANDING
454505
): Promise<PreparedDispatch> {

‎apps/sim/lib/workspace-files/search/enqueue-dispatch.test.ts‎

Lines changed: 34 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,31 +8,31 @@ import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from 'vites
88

99
const mocks = vi.hoisted(() => ({
1010
dispatch: vi.fn(),
11+
hasWork: vi.fn(),
1112
}))
1213

13-
vi.mock('@/background/workspace-file-search-dispatch', () => ({
14-
workspaceFileSearchDispatchTask: {},
15-
}))
1614
vi.mock('@/lib/core/async-jobs/region', () => asyncJobsRegionMock)
1715
vi.mock('@/lib/core/utils/background', () => backgroundTaskMock)
1816
vi.mock('@/lib/workspace-files/search/dispatcher', () => ({
1917
dispatchWorkspaceFileSearchIndexJobs: mocks.dispatch,
18+
hasWorkspaceFileSearchDispatchWork: mocks.hasWork,
2019
}))
2120

2221
import { tasks } from '@trigger.dev/sdk'
2322
import { enqueueWorkspaceFileSearchDispatch } from '@/lib/workspace-files/search/enqueue-dispatch'
2423

2524
const mockTrigger = vi.mocked(tasks.trigger)
2625

27-
setEnvFlags({ isTriggerDevEnabled: true })
2826
afterAll(resetEnvFlagsMock)
2927

3028
describe('workspace file search dispatcher enqueue', () => {
3129
beforeEach(() => {
3230
vi.useFakeTimers()
3331
vi.setSystemTime(new Date('2026-08-29T12:34:45.000Z'))
32+
setEnvFlags({ isTriggerDevEnabled: true })
3433
asyncJobsRegionMockFns.mockResolveTriggerRegion.mockResolvedValue('us-east-1')
3534
mockTrigger.mockResolvedValue({ id: 'run-1' })
35+
mocks.hasWork.mockResolvedValue(true)
3636
})
3737

3838
afterEach(() => {
@@ -41,6 +41,7 @@ describe('workspace file search dispatcher enqueue', () => {
4141

4242
it('waits only for durable Trigger.dev acceptance and does not run the dispatcher inline', async () => {
4343
await expect(enqueueWorkspaceFileSearchDispatch()).resolves.toEqual({
44+
triggered: true,
4445
backend: 'trigger-dev',
4546
jobId: 'run-1',
4647
})
@@ -55,4 +56,33 @@ describe('workspace file search dispatcher enqueue', () => {
5556
expect(mocks.dispatch).not.toHaveBeenCalled()
5657
expect(backgroundTaskMockFns.mockRunDetached).not.toHaveBeenCalled()
5758
})
59+
60+
it.each([
61+
{ backend: 'trigger-dev', triggerDevEnabled: true },
62+
{ backend: 'inline', triggerDevEnabled: false },
63+
])(
64+
'starts no $backend dispatcher run when there is no dispatch work',
65+
async ({ triggerDevEnabled }) => {
66+
setEnvFlags({ isTriggerDevEnabled: triggerDevEnabled })
67+
mocks.hasWork.mockResolvedValue(false)
68+
69+
await expect(enqueueWorkspaceFileSearchDispatch()).resolves.toEqual({
70+
triggered: false,
71+
backend: null,
72+
jobId: null,
73+
})
74+
expect(mockTrigger).not.toHaveBeenCalled()
75+
expect(backgroundTaskMockFns.mockRunDetached).not.toHaveBeenCalled()
76+
}
77+
)
78+
79+
it('still starts a dispatcher run when the work check fails', async () => {
80+
mocks.hasWork.mockRejectedValue(new Error('statement timeout'))
81+
82+
await expect(enqueueWorkspaceFileSearchDispatch()).resolves.toMatchObject({
83+
triggered: true,
84+
backend: 'trigger-dev',
85+
})
86+
expect(mockTrigger).toHaveBeenCalledTimes(1)
87+
})
5888
})

0 commit comments

Comments
 (0)