Skip to content

Commit addb457

Browse files
fix(file-search): recover dispatch claims abandoned between commit and enqueue (#8209)
* fix(file-search): recover dispatch claims abandoned between commit and enqueue * chore(db): format dispatch handoff migration snapshot
1 parent f12677b commit addb457

14 files changed

Lines changed: 29694 additions & 31 deletions

‎apps/sim/lib/workspace-files/search/README.md‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@ The chunk GIN index uses `fastupdate = off`. Each bounded insert updates the mai
1616

1717
Indexing transactions have separate limits from search: ten seconds per statement, five seconds waiting for a lock, and thirty seconds total on PostgreSQL 17. The outer limit leaves time for ordinary statement cancellation and rollback instead of terminating the connection at the same ten-second deadline. PostgreSQL 16 uses the compatible idle-transaction guard. A canceled batch remains unpublished; the task retry starts a fresh fenced build, and cleanup retires the previous attempt. One row's direct GIN insert is not interruptible, so when storage is saturated a single ordinary chunk can run well past the statement deadline and the cancellation lands only after it; smaller batches cannot prevent that. A statement, lock, or transaction timeout is therefore treated as missing database capacity rather than a bad file: the task retries it after about 2, 4, 8, 16, and 30 minutes (with jitter), six attempts in all, so the retries outlast a slow window instead of landing inside it. Other failures keep three attempts with the runner's short default delays. A run waiting to retry still holds one of its workspace's two outstanding dispatch slots. Only a revision that exhausts its attempts is marked failed; this does not automatically retry revisions already marked failed.
1818

19+
A dispatch claim commits before Trigger.dev accepts its run. A dispatcher that stops in between, for example one killed at its 60-second limit while PostgreSQL is still committing, leaves a claim with no run, and that claim holds one of its workspace's two slots. Each claim therefore carries a two-minute handoff deadline in PostgreSQL time, twice the dispatcher's maximum duration. The deadline is cleared once a run is known to exist: the dispatcher clears it after Trigger.dev accepts the batch, and the worker clears it when the build begins. That write skips rows another transaction holds rather than waiting on them, so it cannot deadlock with a bulk file change. The next dispatch releases a claim whose deadline has passed, logs it, and counts it in its result. The file is claimed again later under a new token, which fences out any run the old claim did get, so nothing re-sent has to be deduplicated. A claim with no deadline, including one made before the column existed, keeps the six-hour stale-dispatch recovery, which covers runs lost after they were handed off. In-process dispatch, used when Trigger.dev is not configured, clears the deadline as soon as it schedules the work in its own process, so a restart there still leaves the scheduled claims to the six-hour window.
20+
1921
The indexing task uses an isolated `medium-2x` Trigger worker (4 GB RAM). Document parsers can materialize expanded content before chunking, so source and extracted-text byte limits do not bound parser memory. Parser complexity guards and the worker's memory budget remain separate protections.
2022

2123
File edits, context changes, and deletion invalidate metadata and expire builds. Chunks have no cascading foreign key to files or workspaces. Cleanup locks at most 100 expired builds with `SKIP LOCKED`, deletes at most 1,000 chunks per transaction, retires empty builds in the same batch, and stops after 10 batches or five seconds. Dispatch pauses while at least 10,000 expired chunks await cleanup, so sustained revisions cannot keep admitting new builds faster than retirement can drain them. Existing ready files remain searchable. Stale workers cannot revive a reclaimed build.

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

Lines changed: 16 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@ vi.mock('@/lib/uploads/contexts/workspace', () => ({
2828
getWorkspaceFile: vi.fn(),
2929
fetchWorkspaceFileBuffer: vi.fn(),
3030
}))
31-
vi.mock('@/lib/copilot/tools/server/files/doc-compile', () => ({ resolveServableDoc: vi.fn() }))
31+
vi.mock('@/lib/mothership/tools/server/files/doc-compile', () => ({ resolveServableDoc: vi.fn() }))
3232
vi.mock('@/lib/file-parsers', () => ({ parseBuffer: vi.fn(), isSupportedFileType: vi.fn() }))
3333

3434
import {
@@ -180,6 +180,7 @@ describe('chunked workspace file search on PostgreSQL', () => {
180180
'0358_workspace_file_content_version_precision.sql',
181181
'0359_workspace_file_search_chunks.sql',
182182
ginWriteMigration,
183+
'0382_workspace_file_search_dispatch_handoff.sql',
183184
]) {
184185
await applyMigration(migration)
185186
}
@@ -641,6 +642,20 @@ describe('chunked workspace file search on PostgreSQL', () => {
641642
'pending'
642643
)
643644
})
645+
it('completes a claim handoff when its run begins, never for an older claim', async () => {
646+
const older = new Date('2026-01-01T01:00:00Z')
647+
const newer = new Date('2026-01-01T02:00:00Z')
648+
await connection`UPDATE workspace_file_search_revision
649+
SET dispatched_at = ${newer.toISOString()}::timestamp,
650+
handoff_expires_at = clock_timestamp() + interval '2 minutes'`
651+
const handoff = async () =>
652+
(await connection`SELECT handoff_expires_at FROM workspace_file_search_revision`)[0]
653+
.handoff_expires_at
654+
expect(await beginFileSearchBuild(revision, older.toISOString())).toBeNull()
655+
expect(await handoff()).not.toBeNull()
656+
expect(await beginFileSearchBuild(revision, newer.toISOString())).not.toBeNull()
657+
expect(await handoff()).toBeNull()
658+
})
644659

645660
async function withOccupiedPool(client: postgres.Sql, count: number, run: () => Promise<void>) {
646661
let release!: () => void

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -96,6 +96,14 @@ export const FILE_SEARCH_INDEX_MAX_OUTSTANDING = 100
9696
export const FILE_SEARCH_INDEX_DISPATCH_WORKSPACES = 100
9797
export const FILE_SEARCH_DISPATCH_INTERVAL_MS = 60 * 1000
9898
export const FILE_SEARCH_DISPATCH_MAX_DURATION_SECONDS = 60
99+
/**
100+
* How long a claim may wait for its run to be handed off. A claim commits before Trigger.dev
101+
* accepts the run, so a dispatcher that stops in between leaves a claim with no run. By twice the
102+
* dispatcher task's maximum duration that dispatcher has been stopped, so the next dispatch
103+
* releases the claim instead of waiting out {@link FILE_SEARCH_INDEX_STALE_DISPATCH_MS}. Anything it
104+
* sent that still lands later is fenced out by the claim's token.
105+
*/
106+
export const FILE_SEARCH_DISPATCH_HANDOFF_MS = 2 * FILE_SEARCH_DISPATCH_MAX_DURATION_SECONDS * 1000
99107
/** Leave room for connection setup, rollback, and task failure reporting before the hard cutoff. */
100108
export const FILE_SEARCH_DISPATCH_STATEMENT_TIMEOUT_MS = 10 * 1000
101109
export const FILE_SEARCH_DISPATCH_LOCK_TIMEOUT_MS = 2 * 1000

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

Lines changed: 159 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,11 @@ 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'
28+
import {
29+
FILE_SEARCH_BACKFILL_PAGE_SIZE,
30+
FILE_SEARCH_DISPATCH_HANDOFF_MS,
31+
FILE_SEARCH_INDEX_STALE_DISPATCH_MS,
32+
} from '@/lib/workspace-files/search/constants'
2933
import {
3034
dispatchWorkspaceFileSearchIndexJobs,
3135
prepareWorkspaceFileSearchDispatch,
@@ -80,7 +84,8 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
8084
source_content_updated_at timestamp NOT NULL, status text NOT NULL DEFAULT 'pending',
8185
build_id text, failure_reason text, line_count integer NOT NULL DEFAULT 0,
8286
indexed_bytes integer NOT NULL DEFAULT 0, chunk_count integer NOT NULL DEFAULT 0,
83-
dispatched_at timestamp, updated_at timestamp NOT NULL DEFAULT now()
87+
dispatched_at timestamp, handoff_expires_at timestamp,
88+
updated_at timestamp NOT NULL DEFAULT now()
8489
)`
8590
await connection`CREATE TABLE workspace_file_search_dispatch_queue (
8691
workspace_id text PRIMARY KEY, enqueued_at timestamp NOT NULL,
@@ -297,6 +302,155 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
297302
expect(plan).not.toMatch(/Sort Key: search_index(_\d+)?\.updated_at/)
298303
}, 30_000)
299304

305+
it('releases a claim abandoned before its enqueue once its handoff expires', async () => {
306+
await seedQueue('workspace-1', 3)
307+
/** Preparing without enqueueing is a dispatcher stopped between its commit and its enqueue. */
308+
const abandoned = await prepareWorkspaceFileSearchDispatch()
309+
expect(abandoned.payloads).toHaveLength(2)
310+
const deadlines = await connection`SELECT
311+
extract(epoch FROM handoff_expires_at - clock_timestamp()) * 1000 AS remaining_ms
312+
FROM workspace_file_search_revision WHERE dispatched_at IS NOT NULL`
313+
expect(deadlines).toHaveLength(2)
314+
for (const { remaining_ms } of deadlines) {
315+
expect(Number(remaining_ms)).toBeGreaterThan(FILE_SEARCH_DISPATCH_HANDOFF_MS - 10_000)
316+
expect(Number(remaining_ms)).toBeLessThanOrEqual(FILE_SEARCH_DISPATCH_HANDOFF_MS)
317+
}
318+
/** Until the deadline passes, the claims keep holding their workspace's slots. */
319+
expect(await prepareWorkspaceFileSearchDispatch()).toMatchObject({
320+
payloads: [],
321+
reapedClaims: 0,
322+
abandonedClaims: 0,
323+
})
324+
325+
await connection`UPDATE workspace_file_search_revision
326+
SET handoff_expires_at = clock_timestamp() - interval '1 millisecond'
327+
WHERE dispatched_at IS NOT NULL`
328+
const recovered = await prepareWorkspaceFileSearchDispatch()
329+
expect(recovered.reapedClaims).toBe(2)
330+
expect(recovered.abandonedClaims).toBe(2)
331+
expect(recovered.payloads).toHaveLength(2)
332+
const abandonedToken = abandoned.payloads[0].dispatchToken
333+
if (!abandonedToken) throw new Error('Every claim carries its dispatch token')
334+
expect(recovered.payloads.map((payload) => payload.dispatchToken)).not.toContain(abandonedToken)
335+
const [left] = await connection`SELECT count(*)::int AS claims
336+
FROM workspace_file_search_revision WHERE dispatched_at = ${abandonedToken}::timestamp`
337+
expect(left.claims).toBe(0)
338+
const [unclaimed] = await connection`SELECT count(*)::int AS deadlines
339+
FROM workspace_file_search_revision
340+
WHERE dispatched_at IS NULL AND handoff_expires_at IS NOT NULL`
341+
expect(unclaimed.deadlines).toBe(0)
342+
})
343+
344+
it('leaves an enqueued claim to its run until the stale-dispatch window', async () => {
345+
await seedQueue('workspace-1', 1)
346+
/** A millisecond revision, as file writes store, so the handoff must match it exactly. */
347+
await connection`UPDATE workspace_files SET content_updated_at = '2026-09-16 12:34:56.789'`
348+
await connection`UPDATE workspace_file_search_revision
349+
SET source_content_updated_at = '2026-09-16 12:34:56.789'`
350+
mocks.batchTrigger.mockResolvedValueOnce({ batchId: 'batch-1' })
351+
await expect(dispatchWorkspaceFileSearchIndexJobs()).resolves.toMatchObject({
352+
dispatchedFiles: 1,
353+
})
354+
const [claim] = await connection`SELECT handoff_expires_at FROM workspace_file_search_revision
355+
WHERE dispatched_at IS NOT NULL`
356+
expect(claim.handoff_expires_at).toBeNull()
357+
358+
/** A run can wait in its queue or back off for an hour without being taken for abandoned. */
359+
await connection`UPDATE workspace_file_search_revision
360+
SET dispatched_at = dispatched_at - interval '1 hour' WHERE dispatched_at IS NOT NULL`
361+
expect((await prepareWorkspaceFileSearchDispatch()).reapedClaims).toBe(0)
362+
await connection`UPDATE workspace_file_search_revision
363+
SET dispatched_at = dispatched_at - ${FILE_SEARCH_INDEX_STALE_DISPATCH_MS} * interval '1 millisecond'
364+
WHERE dispatched_at IS NOT NULL`
365+
expect(await prepareWorkspaceFileSearchDispatch()).toMatchObject({
366+
reapedClaims: 1,
367+
abandonedClaims: 0,
368+
})
369+
})
370+
371+
it('does not complete the handoff of a claim released and claimed again meanwhile', async () => {
372+
await seedQueue('workspace-1', 1)
373+
mocks.batchTrigger.mockImplementationOnce(async () => {
374+
/** Another dispatch releases this claim and claims the revision again under its own token. */
375+
await connection`UPDATE workspace_file_search_revision
376+
SET dispatched_at = '2099-01-01', handoff_expires_at = '2099-01-01'
377+
WHERE dispatched_at IS NOT NULL`
378+
return { batchId: 'batch-1' }
379+
})
380+
381+
await dispatchWorkspaceFileSearchIndexJobs()
382+
383+
const [claim] = await connection`SELECT handoff_expires_at::text AS handoff
384+
FROM workspace_file_search_revision WHERE dispatched_at IS NOT NULL`
385+
expect(claim.handoff).toBe('2099-01-01 00:00:00')
386+
})
387+
388+
it('records the handoff around a claim another transaction holds instead of waiting', async () => {
389+
await seedQueue('workspace-1', 2)
390+
let release = () => {}
391+
const held = new Promise<void>((resolve) => {
392+
release = resolve
393+
})
394+
let holder: Promise<unknown> | undefined
395+
mocks.batchTrigger.mockImplementationOnce(async () => {
396+
/** A run beginning its build holds its claim while the handoff is being recorded. */
397+
let locked = () => {}
398+
const lockReady = new Promise<void>((resolve) => {
399+
locked = resolve
400+
})
401+
holder = connection.begin(async (tx) => {
402+
await tx`SELECT file_id FROM workspace_file_search_revision
403+
WHERE file_id = 'workspace-1-000001' FOR UPDATE`
404+
locked()
405+
await held
406+
})
407+
await lockReady
408+
return { batchId: 'batch-1' }
409+
})
410+
try {
411+
await expect(dispatchWorkspaceFileSearchIndexJobs()).resolves.toMatchObject({
412+
dispatchedFiles: 2,
413+
})
414+
const claims = await connection`SELECT file_id,
415+
handoff_expires_at IS NOT NULL AS pending_handoff
416+
FROM workspace_file_search_revision ORDER BY file_id`
417+
expect([...claims]).toEqual([
418+
{ file_id: 'workspace-1-000001', pending_handoff: true },
419+
{ file_id: 'workspace-1-000002', pending_handoff: false },
420+
])
421+
} finally {
422+
release()
423+
await holder
424+
}
425+
})
426+
427+
it('keeps enqueued claims when recording their handoff fails', async () => {
428+
await seedQueue('workspace-1', 1)
429+
await connection`CREATE FUNCTION reject_handoff() RETURNS trigger LANGUAGE plpgsql AS $$
430+
BEGIN
431+
RAISE EXCEPTION 'handoff unavailable';
432+
END
433+
$$`
434+
await connection`CREATE TRIGGER reject_handoff BEFORE UPDATE OF handoff_expires_at
435+
ON workspace_file_search_revision FOR EACH ROW
436+
WHEN (OLD.handoff_expires_at IS NOT NULL AND NEW.handoff_expires_at IS NULL
437+
AND NEW.dispatched_at IS NOT NULL)
438+
EXECUTE FUNCTION reject_handoff()`
439+
try {
440+
mocks.batchTrigger.mockResolvedValueOnce({ batchId: 'batch-1' })
441+
await expect(dispatchWorkspaceFileSearchIndexJobs()).resolves.toMatchObject({
442+
dispatchedFiles: 1,
443+
})
444+
const [claim] = await connection`SELECT dispatched_at, handoff_expires_at
445+
FROM workspace_file_search_revision`
446+
expect(claim.dispatched_at).not.toBeNull()
447+
expect(claim.handoff_expires_at).not.toBeNull()
448+
} finally {
449+
await connection`DROP TRIGGER reject_handoff ON workspace_file_search_revision`
450+
await connection`DROP FUNCTION reject_handoff()`
451+
}
452+
})
453+
300454
it('fails on a locked backfill row and releases the dispatcher lock', async () => {
301455
let release = () => {}
302456
let locked = () => {}
@@ -385,9 +539,11 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
385539
},
386540
}),
387541
])
388-
const [index] = await connection`SELECT dispatched_at FROM workspace_file_search_revision
542+
const [index] =
543+
await connection`SELECT dispatched_at, handoff_expires_at FROM workspace_file_search_revision
389544
WHERE file_id = ${fileId}`
390545
expect(index.dispatched_at).toBeNull()
546+
expect(index.handoff_expires_at).toBeNull()
391547
const [queued] = await connection`SELECT workspace_id FROM workspace_file_search_dispatch_queue
392548
WHERE workspace_id = ${workspaceId}`
393549
expect(queued.workspace_id).toBe(workspaceId)

0 commit comments

Comments
 (0)