Skip to content

Commit 39595a2

Browse files
committed
fix(file-search): isolate and queue search transactions
1 parent 6516b08 commit 39595a2

13 files changed

Lines changed: 343 additions & 42 deletions

File tree

‎apps/sim/lib/workspace-files/application/search-workspace-file-content.ts‎

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -3,15 +3,13 @@ import { loadActiveWorkspaceContext } from '@/lib/uploads/contexts/workspace'
33
import { defineAuthorizedWorkspaceFileUseCase } from '@/lib/workspace-files/application/authorized-workspace-file-use-case'
44
import { fileOperations } from '@/lib/workspace-files/application/operations'
55
import { resolveWorkspaceFolderScope } from '@/lib/workspace-files/resolve-folder-scope'
6+
import { WorkspaceFileSearchUnavailableError } from '@/lib/workspace-files/search/errors'
67
import {
78
compileFileSearchPattern,
89
type FileSearchMode,
910
FileSearchPatternError,
1011
} from '@/lib/workspace-files/search/pattern'
11-
import {
12-
searchWorkspaceFileIndex,
13-
WorkspaceFileSearchUnavailableError,
14-
} from '@/lib/workspace-files/search/repository'
12+
import { searchWorkspaceFileIndex } from '@/lib/workspace-files/search/repository'
1513

1614
export interface SearchWorkspaceFileContentInput {
1715
workspaceId: string

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,9 @@ Search joins the current file revision and resolved workspace/folder scope. A re
2222

2323
For regular chunks, PostgreSQL checks the pattern with newline-aware semantics, then verifies individual logical lines. Long-line fragments use only necessary three-character literals as a conservative prefilter, including all required alternation branches. Two-code-point overlap preserves those literals at every boundary. PostgreSQL reconstructs the complete candidate line and evaluates the original regex, so anchors, word boundaries, repetitions, and arbitrarily long match spans retain line semantics. Fixed overlap alone is never treated as proof of a match. The supported regex grammar and minimum literal requirement are unchanged.
2424

25-
Regular blocks are verified in batches of at most 16 (128 KiB of indexed text); long lines are reconstructed one at a time. Only bounded match-centered previews leave PostgreSQL: at most 201 rows to detect truncation, and at most 2 KiB per rendered result. A single search has a ten-second application deadline with per-statement guards. PostgreSQL 17 additionally enforces a total transaction timeout; PostgreSQL 16 uses the compatible idle-transaction guard. Transaction advisory locks admit at most 20 simultaneous searches per workspace and 5,000 globally per database. The application connection pool queues requests when its connections are occupied. These admission ceilings are operational safeguards, not throughput guarantees. Busy and timed-out searches fail explicitly; they never report an incomplete scan as an authoritative empty result. The reader uses the normal application database connection, so admission is coordinated on the same database as the index.
25+
Regular blocks are verified in batches of at most 16 (128 KiB of indexed text); long lines are reconstructed one at a time. Only bounded match-centered previews leave PostgreSQL: at most 201 rows to detect truncation, and at most 2 KiB per rendered result. A single search has a ten-second application deadline with per-statement guards. PostgreSQL 17 additionally enforces a total transaction timeout; PostgreSQL 16 uses the compatible idle-transaction guard. Transaction advisory locks admit at most 20 simultaneous searches per workspace and 5,000 globally per database. Search transactions use the dedicated `dbFor('search')` primary pool, with five connections per process, so they cannot occupy the application or execution client pools. Before acquiring a connection, a process admits five active searches and at most 100 waiting requests. Queued workspaces rotate after each grant; waiting requests expire after five seconds or leave immediately on cancellation. The local active budget comes from the same pool profile as the driver. A slot is released only when the transaction settles, including errors. The ten-second execution deadline starts after queueing; connecting to the database additionally uses the driver connection timeout. Busy and timed-out searches fail explicitly; they never report an incomplete scan as an authoritative empty result. These bounds apply identically to exact and regex search and do not change result or line-number semantics.
26+
27+
The 20/workspace and 5,000/global advisory ceilings bound admitted transactions across processes; they are not promises of simultaneous execution or throughput. Local queues absorb short bursts without holding database connections. They are not durable jobs or a fleet-wide fair scheduler. The dedicated client pool isolates connection ownership, not PostgreSQL CPU, memory, I/O, or an upstream PgBouncer server pool. Its default URL is the process primary URL; any `DATABASE_URL_SEARCH` override must target the same primary database so current revisions and advisory admission remain coherent. Independent PgBouncer server budgets require separate database/user pool configuration. Total client connections can increase by five per participating process. Before raising execution capacity, measure the number of processes, backend pool budget, queue wait/rejection rates, search latency, and database resource headroom under representative exact and broad-regex workloads. More queueing cannot increase sustained throughput.
2628

2729
Arbitrary regex cannot have a fixed latency guarantee. Common terms, broad alternatives, and punctuation-only literals may require scanning significant scoped text. Larger capacity decisions need representative query plans and workload measurements; neither a per-file byte cap nor a PostgreSQL row-count claim establishes a total corpus capacity.
2830

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,132 @@
1+
/** @vitest-environment node */
2+
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
3+
import { FileSearchAdmission } from '@/lib/workspace-files/search/admission'
4+
import { WorkspaceFileSearchUnavailableError } from '@/lib/workspace-files/search/errors'
5+
6+
describe('FileSearchAdmission', () => {
7+
beforeEach(() => vi.useFakeTimers())
8+
afterEach(() => vi.useRealTimers())
9+
10+
function createAdmission(concurrency = 1, maxPending = 3) {
11+
return new FileSearchAdmission({ concurrency, maxPending, timeoutMs: 5000 })
12+
}
13+
14+
it('queues a burst without exceeding the active budget and releases each slot once', async () => {
15+
const admission = createAdmission(2)
16+
const first = await admission.acquire('a')
17+
const second = await admission.acquire('a')
18+
const granted = vi.fn()
19+
const waiting = admission.acquire('a').then((release) => {
20+
granted()
21+
return release
22+
})
23+
await vi.advanceTimersByTimeAsync(1)
24+
expect(granted).not.toHaveBeenCalled()
25+
first()
26+
const third = await waiting
27+
first()
28+
const fourthGranted = vi.fn()
29+
const fourth = admission.acquire('b').then((release) => {
30+
fourthGranted()
31+
return release
32+
})
33+
await vi.advanceTimersByTimeAsync(1)
34+
expect(fourthGranted).not.toHaveBeenCalled()
35+
second()
36+
;(await fourth)()
37+
third()
38+
expect(vi.getTimerCount()).toBe(0)
39+
})
40+
41+
it('rotates waiting workspaces instead of draining one burst first', async () => {
42+
const admission = createAdmission()
43+
const first = await admission.acquire('a')
44+
const order: string[] = []
45+
const request = (workspace: string) =>
46+
admission.acquire(workspace).then((release) => {
47+
order.push(workspace)
48+
release()
49+
})
50+
const waiting = [request('a'), request('a'), request('b')]
51+
first()
52+
await Promise.all(waiting)
53+
expect(order).toEqual(['a', 'b', 'a'])
54+
})
55+
56+
it('rejects overflow and recovers when the queue drains', async () => {
57+
const admission = createAdmission(1, 1)
58+
const first = await admission.acquire('a')
59+
const waiting = admission.acquire('a')
60+
await expect(admission.acquire('b')).rejects.toBeInstanceOf(WorkspaceFileSearchUnavailableError)
61+
first()
62+
;(await waiting)()
63+
;(await admission.acquire('b'))()
64+
expect(vi.getTimerCount()).toBe(0)
65+
})
66+
67+
it('expires waiting requests without executing them or leaking queue capacity', async () => {
68+
const admission = createAdmission(1, 1)
69+
const first = await admission.acquire('a')
70+
const waiting = expect(admission.acquire('b')).rejects.toBeInstanceOf(
71+
WorkspaceFileSearchUnavailableError
72+
)
73+
await vi.advanceTimersByTimeAsync(5000)
74+
await waiting
75+
const next = admission.acquire('c')
76+
first()
77+
;(await next)()
78+
expect(vi.getTimerCount()).toBe(0)
79+
})
80+
81+
it('checks expiry when granting even if a busy event loop has delayed the timer', async () => {
82+
const admission = createAdmission()
83+
const first = await admission.acquire('a')
84+
const waiting = expect(admission.acquire('b')).rejects.toBeInstanceOf(
85+
WorkspaceFileSearchUnavailableError
86+
)
87+
vi.setSystemTime(Date.now() + 5000)
88+
first()
89+
await waiting
90+
;(await admission.acquire('c'))()
91+
expect(vi.getTimerCount()).toBe(0)
92+
})
93+
94+
it('removes cancelled waiters and their listeners without consuming a connection', async () => {
95+
const admission = createAdmission(1, 1)
96+
const first = await admission.acquire('a')
97+
const controller = new AbortController()
98+
const remove = vi.spyOn(controller.signal, 'removeEventListener')
99+
const reason = new Error('cancelled')
100+
const waiting = expect(admission.acquire('b', controller.signal)).rejects.toBe(reason)
101+
controller.abort(reason)
102+
await waiting
103+
expect(remove).toHaveBeenCalledWith('abort', expect.any(Function))
104+
expect(vi.getTimerCount()).toBe(0)
105+
const next = admission.acquire('c')
106+
first()
107+
;(await next)()
108+
})
109+
110+
it('rejects an already cancelled caller without occupying a slot', async () => {
111+
const admission = createAdmission()
112+
const reason = new Error('cancelled')
113+
await expect(admission.acquire('a', AbortSignal.abort(reason))).rejects.toBe(reason)
114+
;(await admission.acquire('b'))()
115+
})
116+
117+
it('does not recycle an active slot on cancellation until the database work settles', async () => {
118+
const admission = createAdmission()
119+
const controller = new AbortController()
120+
const first = await admission.acquire('a', controller.signal)
121+
controller.abort()
122+
const granted = vi.fn()
123+
const next = admission.acquire('b').then((release) => {
124+
granted()
125+
release()
126+
})
127+
await vi.advanceTimersByTimeAsync(1)
128+
expect(granted).not.toHaveBeenCalled()
129+
first()
130+
await next
131+
})
132+
})
Lines changed: 96 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,96 @@
1+
import { DB_POOL_PROFILES } from '@sim/db/pool-profiles'
2+
import {
3+
FILE_SEARCH_QUEUE_MAX_PENDING,
4+
FILE_SEARCH_QUEUE_TIMEOUT_MS,
5+
} from '@/lib/workspace-files/search/constants'
6+
import { WorkspaceFileSearchUnavailableError } from '@/lib/workspace-files/search/errors'
7+
8+
interface Waiter {
9+
grant: () => void
10+
}
11+
12+
/**
13+
* Bounds search work before it reaches the database pool. Waiting workspaces
14+
* rotate after each grant; cancellation and expiry remove waiters immediately.
15+
*/
16+
export class FileSearchAdmission {
17+
private active = 0
18+
private pending = 0
19+
private readonly workspaces = new Map<string, Set<Waiter>>()
20+
21+
constructor(
22+
private readonly options: { concurrency: number; maxPending: number; timeoutMs: number }
23+
) {}
24+
25+
async acquire(workspaceId: string, signal?: AbortSignal): Promise<() => void> {
26+
signal?.throwIfAborted()
27+
if (this.active < this.options.concurrency) return this.claim()
28+
if (this.pending >= this.options.maxPending) {
29+
throw new WorkspaceFileSearchUnavailableError('Workspace file search is busy. Retry shortly.')
30+
}
31+
32+
return new Promise((resolve, reject) => {
33+
const deadline = Date.now() + this.options.timeoutMs
34+
let settled = false
35+
const remove = () => {
36+
settled = true
37+
clearTimeout(timer)
38+
signal?.removeEventListener('abort', abort)
39+
const waiting = this.workspaces.get(workspaceId)
40+
waiting?.delete(waiter)
41+
if (waiting?.size === 0) this.workspaces.delete(workspaceId)
42+
this.pending--
43+
}
44+
const fail = (error: unknown) => {
45+
if (settled) return
46+
remove()
47+
reject(error)
48+
}
49+
const expire = () =>
50+
fail(
51+
new WorkspaceFileSearchUnavailableError('Workspace file search is busy. Retry shortly.')
52+
)
53+
const abort = () => fail(signal?.reason)
54+
const waiter: Waiter = {
55+
grant: () => {
56+
if (settled) return
57+
if (signal?.aborted) return abort()
58+
if (Date.now() >= deadline) return expire()
59+
remove()
60+
resolve(this.claim())
61+
},
62+
}
63+
const timer = setTimeout(expire, this.options.timeoutMs)
64+
const waiting = this.workspaces.get(workspaceId) ?? new Set<Waiter>()
65+
waiting.add(waiter)
66+
this.workspaces.set(workspaceId, waiting)
67+
this.pending++
68+
signal?.addEventListener('abort', abort, { once: true })
69+
})
70+
}
71+
72+
private claim(): () => void {
73+
this.active++
74+
let released = false
75+
return () => {
76+
if (released) return
77+
released = true
78+
this.active--
79+
while (this.active < this.options.concurrency && this.workspaces.size > 0) {
80+
const [workspaceId, waiters] = this.workspaces.entries().next().value!
81+
const waiter = waiters.values().next().value!
82+
waiter.grant()
83+
if (this.workspaces.has(workspaceId)) {
84+
this.workspaces.delete(workspaceId)
85+
this.workspaces.set(workspaceId, waiters)
86+
}
87+
}
88+
}
89+
}
90+
}
91+
92+
export const fileSearchAdmission = new FileSearchAdmission({
93+
concurrency: DB_POOL_PROFILES.search.primaryMax,
94+
maxPending: FILE_SEARCH_QUEUE_MAX_PENDING,
95+
timeoutMs: FILE_SEARCH_QUEUE_TIMEOUT_MS,
96+
})

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

Lines changed: 71 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { createHash } from 'node:crypto'
22
import { readFileSync, writeFileSync } from 'node:fs'
33
import { resolve } from 'node:path'
4+
import { DB_POOL_PROFILES } from '@sim/db/pool-profiles'
45
import { withUtcTimestamps } from '@sim/db/timestamps'
56
import { generateId } from '@sim/utils/id'
67
import { sql } from 'drizzle-orm'
@@ -9,8 +10,15 @@ import postgres from 'postgres'
910
import { afterAll, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'
1011
import { buildLiteralMatchStart } from '@/lib/workspace-files/search/sql-pattern'
1112

12-
const database = vi.hoisted(() => ({ current: undefined as PostgresJsDatabase | undefined }))
13+
const database = vi.hoisted(() => ({
14+
current: undefined as PostgresJsDatabase | undefined,
15+
search: undefined as PostgresJsDatabase | undefined,
16+
}))
1317
vi.mock('@sim/db', () => ({
18+
dbFor: (role: string) => {
19+
if (role !== 'search' || !database.search) throw new Error('Search database not initialized')
20+
return database.search
21+
},
1422
get db() {
1523
if (!database.current) throw new Error('Test database not initialized')
1624
return database.current
@@ -75,6 +83,16 @@ describe('chunked workspace file search on PostgreSQL', () => {
7583
onnotice: () => {},
7684
})
7785
)
86+
const searchConnection = postgres(
87+
databaseUrl,
88+
withUtcTimestamps({
89+
max: DB_POOL_PROFILES.search.primaryMax,
90+
prepare: false,
91+
fetch_types: false,
92+
connection: { search_path: `${schema},public` },
93+
onnotice: () => {},
94+
})
95+
)
7896

7997
async function addFile(
8098
fileId: string,
@@ -139,7 +157,8 @@ describe('chunked workspace file search on PostgreSQL', () => {
139157
for (const statement of source.split('--> statement-breakpoint'))
140158
if (statement.trim()) await connection.unsafe(statement)
141159
}
142-
database.current = drizzle(connection, {
160+
database.current = drizzle(connection)
161+
database.search = drizzle(searchConnection, {
143162
logger: {
144163
logQuery(query, params) {
145164
if (
@@ -165,7 +184,8 @@ describe('chunked workspace file search on PostgreSQL', () => {
165184
await connection`DROP SCHEMA ${connection(schema)} CASCADE`
166185
} finally {
167186
database.current = undefined
168-
await connection.end()
187+
database.search = undefined
188+
await Promise.all([connection.end(), searchConnection.end()])
169189
}
170190
})
171191

@@ -477,6 +497,54 @@ describe('chunked workspace file search on PostgreSQL', () => {
477497
)
478498
})
479499

500+
async function withOccupiedPool(client: postgres.Sql, count: number, run: () => Promise<void>) {
501+
let release!: () => void
502+
const held = new Promise<void>((resolve) => {
503+
release = resolve
504+
})
505+
let started = 0
506+
const transactions: Promise<unknown>[] = []
507+
const ready = new Promise<void>((resolve, reject) => {
508+
for (let i = 0; i < count; i++) {
509+
transactions.push(
510+
client
511+
.begin(async () => {
512+
if (++started === count) resolve()
513+
await held
514+
})
515+
.catch(reject)
516+
)
517+
}
518+
})
519+
try {
520+
await ready
521+
await run()
522+
} finally {
523+
release()
524+
await Promise.all(transactions)
525+
}
526+
}
527+
528+
it('searches while every shared application connection is occupied', async () => {
529+
await index('needle')
530+
await withOccupiedPool(connection, 4, async () => {
531+
expect((await search('needle')).results).toHaveLength(1)
532+
})
533+
})
534+
535+
it('keeps application queries available while every search connection is occupied', async () => {
536+
await withOccupiedPool(searchConnection, DB_POOL_PROFILES.search.primaryMax, async () => {
537+
expect((await connection`SELECT 1 AS available`)[0].available).toBe(1)
538+
})
539+
})
540+
541+
it('serves a twenty-search workspace burst through the bounded search pool', async () => {
542+
await index('needle')
543+
const results = await Promise.all(Array.from({ length: 20 }, () => search('needle')))
544+
expect(results).toHaveLength(20)
545+
for (const result of results) expect(result.results).toHaveLength(1)
546+
})
547+
480548
it('admits a search up to the workspace and global ceilings', async () => {
481549
const held = await connection.reserve()
482550
try {

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

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,8 @@ export const FILE_SEARCH_CANDIDATE_PAGE_SIZE = 16
2525
export const FILE_SEARCH_CANDIDATE_PROBE_SIZE = 256
2626
export const FILE_SEARCH_QUERY_GLOBAL_CONCURRENCY = 5000
2727
export const FILE_SEARCH_QUERY_WORKSPACE_CONCURRENCY = 20
28+
export const FILE_SEARCH_QUEUE_MAX_PENDING = 100
29+
export const FILE_SEARCH_QUEUE_TIMEOUT_MS = 5000
2830
export const FILE_SEARCH_CANDIDATE_LITERAL_CHARS = 3
2931
export const FILE_SEARCH_BUILD_LEASE_MS = 20 * 60 * 1000
3032
export const FILE_SEARCH_CLEANUP_BATCH_ROWS = 1000
Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,7 @@
1+
/** A retryable search failure caused by contention or index maintenance. */
2+
export class WorkspaceFileSearchUnavailableError extends Error {
3+
constructor(message: string) {
4+
super(message)
5+
this.name = 'WorkspaceFileSearchUnavailableError'
6+
}
7+
}

0 commit comments

Comments
 (0)