Skip to content

Commit d919ee3

Browse files
committed
fix(file-search): bound dispatcher database work
1 parent 25a2136 commit d919ee3

4 files changed

Lines changed: 355 additions & 57 deletions

File tree

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

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,10 @@ export const FILE_SEARCH_INDEX_MAX_OUTSTANDING = 100
4242
export const FILE_SEARCH_INDEX_DISPATCH_WORKSPACES = 100
4343
export const FILE_SEARCH_DISPATCH_INTERVAL_MS = 60 * 1000
4444
export const FILE_SEARCH_DISPATCH_MAX_DURATION_SECONDS = 60
45+
/** Leave room for connection setup, rollback, and task failure reporting before the hard cutoff. */
46+
export const FILE_SEARCH_DISPATCH_STATEMENT_TIMEOUT_MS = 10 * 1000
47+
export const FILE_SEARCH_DISPATCH_LOCK_TIMEOUT_MS = 2 * 1000
48+
export const FILE_SEARCH_DISPATCH_TRANSACTION_TIMEOUT_MS = 20 * 1000
4549
export const FILE_SEARCH_INDEX_MAX_DURATION_SECONDS = 15 * 60
4650
export const FILE_SEARCH_INDEX_STALE_DISPATCH_MS = 6 * 60 * 60 * 1000
4751
export const FILE_SEARCH_INDEX_STALE_REAP_LIMIT = 100
Lines changed: 111 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,111 @@
1+
/** Real PostgreSQL cancellation must roll back preparation and release its advisory lock. */
2+
import { withUtcTimestamps } from '@sim/db/timestamps'
3+
import { getPostgresErrorCode } from '@sim/utils/errors'
4+
import { generateId } from '@sim/utils/id'
5+
import { drizzle, type PostgresJsDatabase } from 'drizzle-orm/postgres-js'
6+
import postgres from 'postgres'
7+
import { afterAll, beforeAll, describe, expect, it, vi } from 'vitest'
8+
9+
const database = vi.hoisted(() => ({ current: undefined as PostgresJsDatabase | undefined }))
10+
vi.mock('@sim/db', () => ({
11+
get db() {
12+
if (!database.current) throw new Error('Dispatcher test database is not initialized')
13+
return database.current
14+
},
15+
}))
16+
vi.mock('@/lib/workspace-files/search/indexing', () => ({
17+
indexWorkspaceFileForSearch: vi.fn(),
18+
markWorkspaceFileSearchIndexFailed: vi.fn(),
19+
}))
20+
21+
import { prepareWorkspaceFileSearchDispatch } from '@/lib/workspace-files/search/dispatcher'
22+
23+
describe('workspace file search dispatch PostgreSQL deadlines', () => {
24+
const schemaName = `dispatch_test_${generateId().replaceAll('-', '')}`
25+
const databaseUrl = process.env.KNOWLEDGE_ACL_TEST_DATABASE_URL
26+
if (!databaseUrl) throw new Error('Dispatcher tests require a disposable local database')
27+
const connection = postgres(
28+
databaseUrl,
29+
withUtcTimestamps({
30+
max: 3,
31+
prepare: false,
32+
fetch_types: false,
33+
connection: { search_path: schemaName },
34+
onnotice: () => {},
35+
})
36+
)
37+
38+
beforeAll(async () => {
39+
await connection`CREATE SCHEMA ${connection(schemaName)}`
40+
await connection`CREATE TABLE workspace_file_search_backfill (
41+
id text PRIMARY KEY, after_workspace_id text, after_file_id text,
42+
completed_at timestamp, updated_at timestamp NOT NULL
43+
)`
44+
await connection`INSERT INTO workspace_file_search_backfill (id, updated_at)
45+
VALUES ('workspace-file-search-v1', '2026-09-16 00:00:00')`
46+
database.current = drizzle(connection)
47+
})
48+
49+
afterAll(async () => {
50+
try {
51+
await connection`DROP SCHEMA ${connection(schemaName)} CASCADE`
52+
} finally {
53+
await connection.end()
54+
database.current = undefined
55+
}
56+
})
57+
58+
async function expectAdvisoryLockReleased() {
59+
await connection.begin(async (tx) => {
60+
const [row] = await tx`SELECT pg_try_advisory_xact_lock(
61+
hashtextextended('workspace-file-search-dispatch', 0)
62+
) AS acquired`
63+
expect(row.acquired).toBe(true)
64+
})
65+
}
66+
67+
it('fails on a locked backfill row and releases the dispatcher lock', async () => {
68+
let release = () => {}
69+
let locked = () => {}
70+
const releaseLock = new Promise<void>((resolve) => {
71+
release = resolve
72+
})
73+
const lockReady = new Promise<void>((resolve) => {
74+
locked = resolve
75+
})
76+
const blocker = connection.begin(async (tx) => {
77+
await tx`SELECT id FROM workspace_file_search_backfill FOR UPDATE`
78+
locked()
79+
await releaseLock
80+
})
81+
await lockReady
82+
try {
83+
const failure = await prepareWorkspaceFileSearchDispatch().catch((error: unknown) => error)
84+
expect(getPostgresErrorCode(failure)).toBe('55P03')
85+
await expectAdvisoryLockReleased()
86+
} finally {
87+
release()
88+
await blocker
89+
}
90+
})
91+
92+
it('cancels a slow statement and rolls back its earlier writes', async () => {
93+
await connection`CREATE FUNCTION slow_backfill() RETURNS trigger LANGUAGE plpgsql AS $$
94+
BEGIN
95+
UPDATE workspace_file_search_backfill SET updated_at = '2099-01-01';
96+
PERFORM pg_sleep(15);
97+
RETURN NEW;
98+
END
99+
$$`
100+
await connection`CREATE TRIGGER slow_backfill BEFORE INSERT ON workspace_file_search_backfill
101+
FOR EACH ROW EXECUTE FUNCTION slow_backfill()`
102+
103+
const failure = await prepareWorkspaceFileSearchDispatch().catch((error: unknown) => error)
104+
105+
expect(getPostgresErrorCode(failure)).toBe('57014')
106+
const [row] =
107+
await connection`SELECT updated_at::text AS updated_at FROM workspace_file_search_backfill`
108+
expect(row.updated_at).toBe('2026-09-16 00:00:00')
109+
await expectAdvisoryLockReleased()
110+
}, 20_000)
111+
})

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

Lines changed: 134 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,43 @@
11
/**
22
* @vitest-environment node
33
*/
4-
import { describe, expect, it } from 'vitest'
4+
import {
5+
workspaceFileSearchBackfill,
6+
workspaceFileSearchDispatchQueue,
7+
workspaceFileSearchIndex,
8+
} from '@sim/db/schema'
9+
import { dbChainMock, dbChainMockFns, queueTableRows, resetDbChainMock } from '@sim/testing'
10+
import { beforeEach, describe, expect, it, vi } from 'vitest'
11+
12+
const mocks = vi.hoisted(() => ({
13+
batchTrigger: vi.fn(),
14+
info: vi.fn(),
15+
error: vi.fn(),
16+
}))
17+
18+
vi.mock('@sim/db/schema', async () => ({
19+
...(await import('@sim/testing/mocks/schema.mock')).schemaMock,
20+
workspaceFileSearchBackfill: { id: 'backfill.id' },
21+
workspaceFileSearchDispatchQueue: {
22+
workspaceId: 'queue.workspaceId',
23+
lastDispatchedAt: 'queue.lastDispatchedAt',
24+
enqueuedAt: 'queue.enqueuedAt',
25+
},
26+
}))
27+
28+
vi.mock('@sim/logger', () => ({ createLogger: () => ({ info: mocks.info, error: mocks.error }) }))
29+
vi.mock('@trigger.dev/sdk', () => ({ tasks: { batchTrigger: mocks.batchTrigger } }))
30+
vi.mock('@/lib/core/config/env-flags', () => ({ isTriggerDevEnabled: true }))
31+
vi.mock('@/lib/core/async-jobs/region', () => ({ resolveTriggerRegion: async () => 'us-east-1' }))
32+
vi.mock('@/lib/workspace-files/search/indexing', () => ({
33+
indexWorkspaceFileForSearch: vi.fn(),
34+
markWorkspaceFileSearchIndexFailed: vi.fn(),
35+
}))
36+
537
import {
638
buildWorkspaceFileSearchTriggerItems,
39+
dispatchWorkspaceFileSearchIndexJobs,
40+
prepareWorkspaceFileSearchDispatch,
741
shouldUseWorkspaceFileSearchTrigger,
842
} from '@/lib/workspace-files/search/dispatcher'
943

@@ -34,3 +68,102 @@ describe('workspace file search dispatch policy', () => {
3468
])
3569
})
3670
})
71+
72+
describe('workspace file search dispatch deadlines', () => {
73+
beforeEach(() => {
74+
vi.clearAllMocks()
75+
resetDbChainMock()
76+
})
77+
78+
it('sets local database deadlines before taking the advisory lock', async () => {
79+
dbChainMockFns.execute.mockResolvedValueOnce([]).mockResolvedValueOnce([{ acquired: false }])
80+
81+
await expect(prepareWorkspaceFileSearchDispatch()).resolves.toEqual({
82+
payloads: [],
83+
backfilledFiles: 0,
84+
reapedClaims: 0,
85+
lockAcquired: false,
86+
})
87+
88+
const guards = JSON.stringify(dbChainMockFns.execute.mock.calls[0][0])
89+
expect(guards).toContain("set_config('statement_timeout', ")
90+
expect(guards).toContain('10000ms')
91+
expect(guards).toContain("set_config('lock_timeout', ")
92+
expect(guards).toContain('2000ms')
93+
expect(guards).toContain("set_config('transaction_timeout', ")
94+
expect(guards).toContain('20000ms')
95+
expect(guards.match(/, true\)/g)).toHaveLength(3)
96+
expect(JSON.stringify(dbChainMockFns.execute.mock.calls[1][0])).toContain(
97+
'pg_try_advisory_xact_lock'
98+
)
99+
expect(dbChainMockFns.insert).not.toHaveBeenCalled()
100+
})
101+
102+
it.each(['57014', '55P03', '25P04'])(
103+
'propagates SQLSTATE %s without enqueuing an uncommitted claim',
104+
async (code) => {
105+
const error = new Error('Failed query\nparams: sensitive-value', {
106+
cause: Object.assign(new Error('database timeout'), { code }),
107+
})
108+
dbChainMockFns.execute.mockResolvedValueOnce([]).mockResolvedValueOnce([{ acquired: true }])
109+
dbChainMockFns.onConflictDoNothing.mockRejectedValueOnce(error)
110+
111+
await expect(dispatchWorkspaceFileSearchIndexJobs()).rejects.toBe(error)
112+
113+
expect(mocks.batchTrigger).not.toHaveBeenCalled()
114+
expect(mocks.error).toHaveBeenCalledWith('Workspace file search dispatch phase failed', {
115+
phase: 'backfill',
116+
durationMs: expect.any(Number),
117+
code,
118+
error: 'Failed query',
119+
})
120+
expect(JSON.stringify(mocks.error.mock.calls)).not.toContain('sensitive-value')
121+
}
122+
)
123+
124+
it('reports a transaction failure even after the transaction callback finishes', async () => {
125+
const error = Object.assign(new Error('commit failed'), { code: '08006' })
126+
dbChainMockFns.execute.mockResolvedValueOnce([]).mockResolvedValueOnce([{ acquired: false }])
127+
dbChainMockFns.transaction.mockImplementationOnce(async (callback) => {
128+
await callback(dbChainMock.db)
129+
throw error
130+
})
131+
132+
await expect(prepareWorkspaceFileSearchDispatch()).rejects.toBe(error)
133+
134+
expect(mocks.error).toHaveBeenCalledWith('Workspace file search dispatch phase failed', {
135+
phase: 'prepare-transaction',
136+
durationMs: expect.any(Number),
137+
code: '08006',
138+
error: 'commit failed',
139+
})
140+
})
141+
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+
})
169+
})

0 commit comments

Comments
 (0)