Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions apps/sim/app/api/cron/workspace-file-search-dispatch/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
113 changes: 107 additions & 6 deletions apps/sim/lib/workspace-files/search/dispatcher.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Comment on lines +21 to 23

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 The mock loads real imports

The new partial mock uses importOriginal. The testing guide explicitly says never to use vi.importActual() or importOriginal to build a partial mock; use the central mock instead. This loads the real module and its imports just to retain fileSearchBuildExpired. 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!

}))
vi.mock('@trigger.dev/sdk', () => ({ tasks: { batchTrigger: mocks.batchTrigger } }))
Expand All @@ -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'

Expand Down Expand Up @@ -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`
})

Expand Down Expand Up @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 The test checks another index

The new plan test expects a fixture index on (workspace_id, dispatched_at), but production's workspace_file_search_revision_active_idx uses (dispatched_at, workspace_id). It therefore does not check the production index promised by the test. Create that index in the fixture and assert its name. Also distinguish forced index use from the plan PostgreSQL would normally choose.

})
})
})
87 changes: 69 additions & 18 deletions apps/sim/lib/workspace-files/search/dispatcher.ts
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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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.
*
Expand All @@ -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

Expand Down Expand Up @@ -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,
Expand All @@ -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 })
Expand Down Expand Up @@ -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

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 The check can stall dispatch

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, hasDispatchWork never reaches its catch and no dispatcher starts. Run the check in a short, timeout-bound transaction so a stalled check can fall back to dispatching.

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> {
Expand Down
38 changes: 34 additions & 4 deletions apps/sim/lib/workspace-files/search/enqueue-dispatch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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(() => {
Expand All @@ -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',
})
Expand All @@ -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)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 The test checks a spy

The new failure-path test proves dispatch by asserting that mockTrigger was called. CLAUDE.md explicitly forbids tests that assert mock calls. Check the actual enqueue boundary and its accepted result instead, so the test catches a failed handoff rather than just a call to a spy. 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!

})
})
Loading
Loading