Skip to content

Commit 20a36ad

Browse files
committed
fix(file-search): preserve compatible dispatch cleanup
1 parent d919ee3 commit 20a36ad

4 files changed

Lines changed: 170 additions & 38 deletions

File tree

‎.github/workflows/test-build.yml‎

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -134,6 +134,23 @@ jobs:
134134
BILLING_USAGE_TEST_DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5433/sim_billing_test
135135
run: bunx vitest run lib/billing/core/usage-log.postgres.test.ts
136136

137+
- name: Verify file search dispatch deadlines on PostgreSQL 17
138+
working-directory: apps/sim
139+
env:
140+
TZ: America/Los_Angeles
141+
KNOWLEDGE_ACL_TEST_DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5432/sim_auth_scim
142+
run: bunx vitest run --mode integration lib/workspace-files/search/dispatcher.integration.ts
143+
144+
- name: Verify file search dispatch deadlines on PostgreSQL 16
145+
if: matrix.provision == 'push'
146+
working-directory: apps/sim
147+
env:
148+
TZ: America/Los_Angeles
149+
KNOWLEDGE_ACL_TEST_DATABASE_URL: postgresql://postgres:postgres@127.0.0.1:5433/sim_acl_test_dispatch
150+
run: |
151+
bun -e 'import postgres from "postgres"; const url = new URL(process.env.KNOWLEDGE_ACL_TEST_DATABASE_URL); const database = url.pathname.slice(1); url.pathname = "/postgres"; const sql = postgres(url.toString()); await sql`CREATE DATABASE ${sql(database)}`; await sql.end()'
152+
bunx vitest run --mode integration lib/workspace-files/search/dispatcher.integration.ts
153+
137154
- name: Verify SCIM and administration over real HTTP
138155
working-directory: apps/sim
139156
env:

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

Lines changed: 82 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,12 +1,14 @@
11
/** Real PostgreSQL cancellation must roll back preparation and release its advisory lock. */
22
import { withUtcTimestamps } from '@sim/db/timestamps'
33
import { getPostgresErrorCode } from '@sim/utils/errors'
4+
import { sleep } from '@sim/utils/helpers'
45
import { generateId } from '@sim/utils/id'
56
import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js'
67
import postgres from 'postgres'
7-
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
8+
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
89

910
const database = vi.hoisted(() => ({ current: undefined as PostgresJsDatabase | undefined }))
11+
const mocks = vi.hoisted(() => ({ batchTrigger: vi.fn() }))
1012
vi.mock('@sim/db', () => ({
1113
get db() {
1214
if (!database.current) throw new Error('Dispatcher test database is not initialized')
@@ -17,8 +19,14 @@ vi.mock('@/lib/workspace-files/search/indexing', () => ({
1719
indexWorkspaceFileForSearch: vi.fn(),
1820
markWorkspaceFileSearchIndexFailed: vi.fn(),
1921
}))
22+
vi.mock('@trigger.dev/sdk', () => ({ tasks: { batchTrigger: mocks.batchTrigger } }))
23+
vi.mock('@/lib/core/config/env-flags', () => ({ isTriggerDevEnabled: true }))
24+
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
2025

21-
import { prepareWorkspaceFileSearchDispatch } from '@/lib/workspace-files/search/dispatcher'
26+
import {
27+
dispatchWorkspaceFileSearchIndexJobs,
28+
prepareWorkspaceFileSearchDispatch,
29+
} from '@/lib/workspace-files/search/dispatcher'
2230

2331
describe('workspace file search dispatch PostgreSQL deadlines', () => {
2432
const schemaName = `dispatch_test_${generateId().replaceAll('-', '')}`
@@ -41,11 +49,32 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
4149
id text PRIMARY KEY, after_workspace_id text, after_file_id text,
4250
completed_at timestamp, updated_at timestamp NOT NULL
4351
)`
52+
await connection`CREATE TABLE workspace_files (
53+
id text PRIMARY KEY, workspace_id text NOT NULL, context text NOT NULL,
54+
deleted_at timestamp, content_updated_at timestamp NOT NULL
55+
)`
56+
await connection`CREATE TABLE workspace_file_search_index (
57+
file_id text NOT NULL, workspace_id text NOT NULL, source_content_updated_at timestamp NOT NULL,
58+
status text NOT NULL, dispatched_at timestamp, updated_at timestamp NOT NULL,
59+
PRIMARY KEY (file_id, source_content_updated_at)
60+
)`
61+
await connection`CREATE TABLE workspace_file_search_dispatch_queue (
62+
workspace_id text PRIMARY KEY, enqueued_at timestamp NOT NULL,
63+
updated_at timestamp NOT NULL, last_dispatched_at timestamp
64+
)`
4465
await connection`INSERT INTO workspace_file_search_backfill (id, updated_at)
4566
VALUES ('workspace-file-search-v1', '2026-09-16 00:00:00')`
4667
database.current = drizzle(connection)
4768
})
4869

70+
beforeEach(async () => {
71+
mocks.batchTrigger.mockReset()
72+
await connection`DROP TRIGGER IF EXISTS slow_backfill ON workspace_file_search_backfill`
73+
await connection`TRUNCATE workspace_files, workspace_file_search_index, workspace_file_search_dispatch_queue`
74+
await connection`UPDATE workspace_file_search_backfill
75+
SET updated_at = '2026-09-16 00:00:00', completed_at = NULL`
76+
})
77+
4978
afterAll(async () => {
5079
try {
5180
await connection`DROP SCHEMA ${connection(schemaName)} CASCADE`
@@ -108,4 +137,55 @@ describe('workspace file search dispatch PostgreSQL deadlines', () => {
108137
expect(row.updated_at).toBe('2026-09-16 00:00:00')
109138
await expectAdvisoryLockReleased()
110139
}, 20_000)
140+
141+
it('releases committed claims after contention exceeds the preparation lock deadline', async () => {
142+
const fileId = generateId()
143+
const workspaceId = generateId()
144+
await connection`UPDATE workspace_file_search_backfill SET completed_at = now()`
145+
await connection`INSERT INTO workspace_files (id, workspace_id, context, content_updated_at)
146+
VALUES (${fileId}, ${workspaceId}, 'workspace', '2026-09-16')`
147+
await connection`INSERT INTO workspace_file_search_index
148+
(file_id, workspace_id, source_content_updated_at, status, updated_at)
149+
VALUES (${fileId}, ${workspaceId}, '2026-09-16', 'pending', now())`
150+
await connection`INSERT INTO workspace_file_search_dispatch_queue
151+
(workspace_id, enqueued_at, updated_at) VALUES (${workspaceId}, now(), now())`
152+
153+
const enqueueError = new Error('Queue unavailable')
154+
let blocker: Promise<unknown> | undefined
155+
mocks.batchTrigger.mockImplementationOnce(async () => {
156+
let locked = () => {}
157+
const lockReady = new Promise<void>((resolve) => {
158+
locked = resolve
159+
})
160+
blocker = connection.begin(async (tx) => {
161+
await tx`SELECT file_id FROM workspace_file_search_index WHERE file_id = ${fileId} FOR UPDATE`
162+
locked()
163+
await sleep(3_000)
164+
})
165+
await Promise.race([lockReady, blocker])
166+
throw enqueueError
167+
})
168+
169+
try {
170+
await expect(dispatchWorkspaceFileSearchIndexJobs()).rejects.toBe(enqueueError)
171+
expect(mocks.batchTrigger).toHaveBeenCalledWith('workspace-file-search-index', [
172+
expect.objectContaining({
173+
payload: {
174+
fileId,
175+
workspaceId,
176+
sourceContentUpdatedAt: '2026-09-16T00:00:00.000Z',
177+
},
178+
}),
179+
])
180+
const [index] = await connection`SELECT dispatched_at FROM workspace_file_search_index
181+
WHERE file_id = ${fileId}`
182+
expect(index.dispatched_at).toBeNull()
183+
const [queued] =
184+
await connection`SELECT workspace_id FROM workspace_file_search_dispatch_queue
185+
WHERE workspace_id = ${workspaceId}`
186+
expect(queued.workspace_id).toBe(workspaceId)
187+
} finally {
188+
await blocker
189+
}
190+
})
111191
})

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

Lines changed: 48 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -90,9 +90,10 @@ describe('workspace file search dispatch deadlines', () => {
9090
expect(guards).toContain('10000ms')
9191
expect(guards).toContain("set_config('lock_timeout', ")
9292
expect(guards).toContain('2000ms')
93-
expect(guards).toContain("set_config('transaction_timeout', ")
93+
expect(guards).toContain("current_setting('transaction_timeout', true) is null")
94+
expect(guards).toContain("then 'idle_in_transaction_session_timeout'")
95+
expect(guards).toContain("else 'transaction_timeout'")
9496
expect(guards).toContain('20000ms')
95-
expect(guards.match(/, true\)/g)).toHaveLength(3)
9697
expect(JSON.stringify(dbChainMockFns.execute.mock.calls[1][0])).toContain(
9798
'pg_try_advisory_xact_lock'
9899
)
@@ -139,31 +140,49 @@ describe('workspace file search dispatch deadlines', () => {
139140
})
140141
})
141142

142-
it('also bounds releasing committed claims when batch submission fails', async () => {
143-
queueTableRows(workspaceFileSearchBackfill, [{ completedAt: new Date() }])
144-
queueTableRows(workspaceFileSearchIndex, [])
145-
queueTableRows(workspaceFileSearchIndex, [{ active: 0 }])
146-
queueTableRows(workspaceFileSearchDispatchQueue, [{ workspaceId: 'workspace-1' }])
147-
dbChainMockFns.execute
148-
.mockResolvedValueOnce([])
149-
.mockResolvedValueOnce([{ acquired: true }])
150-
.mockResolvedValueOnce([
151-
{
152-
workspaceId: 'workspace-1',
153-
fileId: 'file-1',
154-
sourceContentUpdatedAt: new Date('2026-09-16T00:00:00Z'),
155-
},
156-
])
157-
const error = new Error('Trigger unavailable')
158-
mocks.batchTrigger.mockRejectedValueOnce(error)
159-
160-
await expect(dispatchWorkspaceFileSearchIndexJobs()).rejects.toBe(error)
161-
162-
expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(2)
163-
const guards = dbChainMockFns.execute.mock.calls.filter(([query]) =>
164-
JSON.stringify(query).includes('statement_timeout')
165-
)
166-
expect(guards).toHaveLength(2)
167-
expect(dbChainMockFns.set).toHaveBeenCalledWith(expect.objectContaining({ dispatchedAt: null }))
168-
})
143+
it.each([false, true])(
144+
'preserves enqueue failures when claim release fails: %s',
145+
async (releaseFails) => {
146+
queueTableRows(workspaceFileSearchBackfill, [{ completedAt: new Date() }])
147+
queueTableRows(workspaceFileSearchIndex, [])
148+
queueTableRows(workspaceFileSearchIndex, [{ active: 0 }])
149+
queueTableRows(workspaceFileSearchDispatchQueue, [{ workspaceId: 'workspace-1' }])
150+
dbChainMockFns.execute
151+
.mockResolvedValueOnce([])
152+
.mockResolvedValueOnce([{ acquired: true }])
153+
.mockResolvedValueOnce([
154+
{
155+
workspaceId: 'workspace-1',
156+
fileId: 'file-1',
157+
sourceContentUpdatedAt: new Date('2026-09-16T00:00:00Z'),
158+
},
159+
])
160+
const error = new Error('Trigger unavailable')
161+
const releaseError = new Error('claim release unavailable')
162+
mocks.batchTrigger.mockRejectedValueOnce(error)
163+
if (releaseFails) {
164+
dbChainMockFns.transaction
165+
.mockImplementationOnce(async (callback) => callback(dbChainMock.db))
166+
.mockRejectedValueOnce(releaseError)
167+
}
168+
169+
if (releaseFails) {
170+
await expect(dispatchWorkspaceFileSearchIndexJobs()).rejects.toMatchObject({
171+
errors: [error, releaseError],
172+
cause: error,
173+
})
174+
} else {
175+
await expect(dispatchWorkspaceFileSearchIndexJobs()).rejects.toBe(error)
176+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
177+
expect.objectContaining({ dispatchedAt: null })
178+
)
179+
}
180+
181+
expect(dbChainMockFns.transaction).toHaveBeenCalledTimes(2)
182+
const guards = dbChainMockFns.execute.mock.calls.filter(([query]) =>
183+
JSON.stringify(query).includes('statement_timeout')
184+
)
185+
expect(guards).toHaveLength(1)
186+
}
187+
)
169188
})

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

Lines changed: 23 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -72,13 +72,20 @@ async function runDispatchPhase<T>(phase: string, operation: () => Promise<T>):
7272
}
7373
}
7474

75-
/** Transaction-local guards also cover claim release; a stalled query must roll back before retry. */
75+
/** Preparation must roll back before retry; PostgreSQL 16 falls back to an idle transaction guard. */
7676
async function configureDispatchTimeouts(tx: DbTransaction): Promise<void> {
7777
await tx.execute(sql`
7878
SELECT
7979
set_config('statement_timeout', ${`${FILE_SEARCH_DISPATCH_STATEMENT_TIMEOUT_MS}ms`}, true),
8080
set_config('lock_timeout', ${`${FILE_SEARCH_DISPATCH_LOCK_TIMEOUT_MS}ms`}, true),
81-
set_config('transaction_timeout', ${`${FILE_SEARCH_DISPATCH_TRANSACTION_TIMEOUT_MS}ms`}, true)
81+
set_config(
82+
case when current_setting('transaction_timeout', true) is null
83+
then 'idle_in_transaction_session_timeout'
84+
else 'transaction_timeout'
85+
end,
86+
${`${FILE_SEARCH_DISPATCH_TRANSACTION_TIMEOUT_MS}ms`},
87+
true
88+
)
8289
`)
8390
}
8491

@@ -311,7 +318,7 @@ async function claimQueuedWorkspaceJobs(
311318
const rows = await tx.execute<{
312319
workspaceId: string
313320
fileId: string
314-
sourceContentUpdatedAt: Date
321+
sourceContentUpdatedAt: string
315322
}>(sql`
316323
WITH selected_workspace(workspace_id) AS (
317324
VALUES ${workspaceValues}
@@ -370,7 +377,7 @@ async function claimQueuedWorkspaceJobs(
370377
RETURNING
371378
search_index.workspace_id AS "workspaceId",
372379
search_index.file_id AS "fileId",
373-
search_index.source_content_updated_at AS "sourceContentUpdatedAt"
380+
search_index.source_content_updated_at AT TIME ZONE 'UTC' AS "sourceContentUpdatedAt"
374381
`)
375382

376383
const remainingForWorkspace = tx
@@ -484,7 +491,6 @@ async function releaseDispatchClaims(payloads: readonly WorkspaceFileSearchIndex
484491
}))
485492
await runDispatchPhase('release-claims', () =>
486493
db.transaction(async (tx) => {
487-
await configureDispatchTimeouts(tx)
488494
const filter = revisionFilter(rows)
489495
if (filter) {
490496
await tx
@@ -555,11 +561,21 @@ export async function dispatchWorkspaceFileSearchIndexJobs(): Promise<WorkspaceF
555561
lockAcquired: prepared.lockAcquired,
556562
}
557563
} catch (error) {
558-
await releaseDispatchClaims(prepared.payloads)
559564
logger.error('Failed to dispatch workspace file search indexing batch', {
560565
files: prepared.payloads.length,
561-
error: getErrorMessage(error),
566+
error: truncate(getErrorMessage(error).split('\nparams: ')[0], 500),
562567
})
568+
try {
569+
await releaseDispatchClaims(prepared.payloads)
570+
} catch (releaseError) {
571+
throw new AggregateError(
572+
[error, releaseError],
573+
'File search enqueue and claim release failed',
574+
{
575+
cause: error,
576+
}
577+
)
578+
}
563579
throw error
564580
}
565581
}

0 commit comments

Comments
 (0)