Skip to content

Commit 0ae1961

Browse files
committed
fix(search): harden live search accuracy, provider queries, and request cost
- GitHub: batch code search by bytes under its 1,000-byte query limit, send date bounds as one updated:start..end range, keep qualifiers outside grouped text, and exclude dateless code from dated searches - Calendar: keep one verified copy of a meeting shared across calendars and order dated agendas by start - Clean Gmail, Calendar, and Confluence text; report Google quota 403s as rate limits - Report partial only for degraded coverage; drop cursors that would skip merged-out results - Share one account session between search and read; filter before verifying, verify in parallel, pass GitLab evidence, and reuse listed accounts - Pin DNS on the request as well as the agent (Bun ignores agent lookups), and reuse pinned keep-alive connections with compressed responses - Route the Search MCP chat tool to the assistant chat endpoint
1 parent 2d0353a commit 0ae1961

35 files changed

Lines changed: 1837 additions & 688 deletions

‎apps/sim/connectors/google-workspace/api-errors.ts‎

Lines changed: 8 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,13 @@ const RATE_LIMIT_REASONS = new Set([
5454
'RATE_LIMIT_EXCEEDED',
5555
])
5656

57+
/** Whether a Google error reason reports an exhausted rate or usage quota rather than a denial. */
58+
export function isGoogleQuotaReason(reason: string): boolean {
59+
return (
60+
RATE_LIMIT_REASONS.has(reason) || reason === 'dailyLimitExceeded' || reason === 'quotaExceeded'
61+
)
62+
}
63+
5764
export function safeGoogleErrorReasons(reasons: readonly string[]): string[] {
5865
return [...new Set(reasons.filter((reason) => SAFE_REASONS.has(reason)))]
5966
}
@@ -119,13 +126,7 @@ export class GoogleApiError extends ConnectorSourceError {
119126
const safeReasons = safeGoogleErrorReasons(reasons)
120127
const suffix = safeReasons.length ? ` (${safeReasons.join(', ')})` : ''
121128
const category =
122-
status === 429 ||
123-
safeReasons.some(
124-
(reason) =>
125-
RATE_LIMIT_REASONS.has(reason) ||
126-
reason === 'dailyLimitExceeded' ||
127-
reason === 'quotaExceeded'
128-
)
129+
status === 429 || safeReasons.some(isGoogleQuotaReason)
129130
? 'rate_limit'
130131
: status >= 500
131132
? 'provider_unavailable'

‎apps/sim/connectors/utils.ts‎

Lines changed: 15 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -362,15 +362,25 @@ export function looksLikeHtml(value: string): boolean {
362362
* punctuation as numeric references, which previously reached the index verbatim.
363363
*/
364364
export function htmlToPlainText(html: string): string {
365-
const text = html
366-
.replace(/<[^>]*>/g, ' ')
367-
.replace(HTML_ENTITY_PATTERN, (raw: string, hex?: string, decimal?: string, named?: string) => {
365+
return decodeHtmlEntities(html.replace(/<[^>]*>/g, ' '))
366+
.replace(/\s+/g, ' ')
367+
.trim()
368+
}
369+
370+
/**
371+
* Decodes HTML character references without touching markup or whitespace. Use for text a
372+
* provider HTML-escapes but does not mark up, such as Gmail message snippets.
373+
*/
374+
export function decodeHtmlEntities(text: string): string {
375+
return text.replace(
376+
HTML_ENTITY_PATTERN,
377+
(raw: string, hex?: string, decimal?: string, named?: string) => {
368378
if (named !== undefined) return NAMED_ENTITIES[named] ?? raw
369379
if (hex !== undefined) return decodeCharacterReference(raw, Number.parseInt(hex, 16))
370380
if (decimal !== undefined) return decodeCharacterReference(raw, Number.parseInt(decimal, 10))
371381
return raw
372-
})
373-
return text.replace(/\s+/g, ' ').trim()
382+
}
383+
)
374384
}
375385

376386
/**

‎apps/sim/lib/core/security/input-validation.server.ts‎

Lines changed: 61 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -356,6 +356,17 @@ export interface SecureFetchOptions {
356356
proxyUrl?: string
357357
/** Hide credential-derived URL details from validation logs. */
358358
logUrlValidationDetails?: boolean
359+
/**
360+
* Ask for a gzip, deflate, or brotli body. The body is decoded before it is returned, and
361+
* `maxResponseBytes` bounds the decoded bytes, so a compression bomb still stops at the cap.
362+
*/
363+
acceptCompressed?: boolean
364+
/**
365+
* Reuses keep-alive connections to the same pinned address across requests. A connection is
366+
* only ever reused for the IP it was opened to, so every request keeps its DNS pinning. The
367+
* owner must call {@link PinnedConnectionPool.destroy} once its requests have finished.
368+
*/
369+
connectionPool?: PinnedConnectionPool
359370
/**
360371
* Where this request's URL came from. Carried on the options so the same
361372
* policy is re-applied to every redirect hop rather than re-derived — a hop
@@ -447,6 +458,43 @@ function resolveRedirectUrl(baseUrl: string, location: string): string {
447458
}
448459
}
449460

461+
/** Keep-alive agents keyed by protocol, host, port, and the pinned address they connect to. */
462+
export interface PinnedConnectionPool {
463+
/** Undefined once destroyed, so a late request falls back to a single-use pinned agent. */
464+
agent(isHttps: boolean, host: string, port: number, resolvedIP: string): http.Agent | undefined
465+
destroy(): void
466+
}
467+
468+
/**
469+
* Creates a request-scoped pool of pinned keep-alive agents. Reusing a connection skips the TCP
470+
* and TLS handshakes that otherwise dominate short provider API calls.
471+
*/
472+
export function createPinnedConnectionPool(): PinnedConnectionPool {
473+
const agents = new Map<string, http.Agent>()
474+
let destroyed = false
475+
return {
476+
agent(isHttps, host, port, resolvedIP) {
477+
if (destroyed) return undefined
478+
const key = JSON.stringify([isHttps, host, port, resolvedIP])
479+
let agent = agents.get(key)
480+
if (!agent) {
481+
const options: http.AgentOptions = {
482+
keepAlive: true,
483+
lookup: createPinnedLookup(resolvedIP),
484+
}
485+
agent = isHttps ? new https.Agent(options) : new http.Agent(options)
486+
agents.set(key, agent)
487+
}
488+
return agent
489+
},
490+
destroy() {
491+
destroyed = true
492+
for (const agent of agents.values()) agent.destroy()
493+
agents.clear()
494+
},
495+
}
496+
}
497+
450498
/**
451499
* Creates a DNS lookup function that always returns a pre-resolved IP address.
452500
* Use this to prevent DNS rebinding (TOCTOU) attacks when connecting to
@@ -1088,6 +1136,11 @@ export async function secureFetchWithPinnedIP(
10881136
const port = parsed.port ? Number.parseInt(parsed.port, 10) : defaultPort
10891137

10901138
let agent: http.Agent | undefined
1139+
/**
1140+
* Bun ignores a `lookup` set on an Agent and honors one on the request, while Node honors
1141+
* both. A pinned direct connection sets it in both places so pinning holds in either runtime.
1142+
*/
1143+
let pinnedLookup: LookupFunction | undefined
10911144
if (outboundDispatcher) {
10921145
agent = undefined
10931146
} else if (options.proxyUrl) {
@@ -1096,12 +1149,16 @@ export async function secureFetchWithPinnedIP(
10961149
// targets tunnel via CONNECT, http targets use absolute-URI forwarding.
10971150
agent = isHttps ? new HttpsProxyAgent(options.proxyUrl) : new HttpProxyAgent(options.proxyUrl)
10981151
} else {
1099-
const lookup = createPinnedLookup(resolvedIP)
1100-
const agentOptions: http.AgentOptions = { lookup }
1101-
agent = isHttps ? new https.Agent(agentOptions) : new http.Agent(agentOptions)
1152+
pinnedLookup = createPinnedLookup(resolvedIP)
1153+
agent =
1154+
options.connectionPool?.agent(isHttps, parsed.hostname, port, resolvedIP) ??
1155+
(isHttps
1156+
? new https.Agent({ lookup: pinnedLookup })
1157+
: new http.Agent({ lookup: pinnedLookup }))
11021158
}
11031159

11041160
const { 'accept-encoding': _, ...sanitizedHeaders } = options.headers ?? {}
1161+
if (options.acceptCompressed) sanitizedHeaders['accept-encoding'] = 'gzip, deflate, br'
11051162
if (!Object.keys(sanitizedHeaders).some((name) => name.toLowerCase() === 'user-agent')) {
11061163
sanitizedHeaders['user-agent'] = DEFAULT_USER_AGENT
11071164
}
@@ -1123,6 +1180,7 @@ export async function secureFetchWithPinnedIP(
11231180
method: options.method || 'GET',
11241181
headers: sanitizedHeaders,
11251182
agent,
1183+
...(pinnedLookup ? { lookup: pinnedLookup } : {}),
11261184
timeout: options.timeout || 300000,
11271185
}
11281186

Lines changed: 148 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,148 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import http from 'node:http'
5+
import type { AddressInfo } from 'node:net'
6+
import { gzipSync } from 'node:zlib'
7+
import { afterEach, describe, expect, it, vi } from 'vitest'
8+
9+
vi.mock('@sim/security/dns', () => ({
10+
resolveHostAddresses: vi.fn(),
11+
preferIpv4: (addresses: string[]) => addresses[0],
12+
}))
13+
14+
vi.mock('@/lib/core/config/env-flags', () => ({
15+
isHosted: false,
16+
getEgressAllowedHosts: () => undefined,
17+
getEgressAllowedIpRanges: () => undefined,
18+
isLegacyPrivateDatabaseAccessAllowed: () => false,
19+
getProxyUrl: () => undefined,
20+
}))
21+
22+
import {
23+
createPinnedConnectionPool,
24+
secureFetchWithPinnedIP,
25+
} from '@/lib/core/security/input-validation.server'
26+
27+
const servers: http.Server[] = []
28+
29+
afterEach(() => {
30+
for (const server of servers.splice(0)) {
31+
server.closeAllConnections()
32+
server.close()
33+
}
34+
})
35+
36+
/** Starts a loopback server that counts the TCP connections it accepts. */
37+
async function startServer(handler: http.RequestListener) {
38+
const server = http.createServer(handler)
39+
servers.push(server)
40+
let connections = 0
41+
server.on('connection', () => connections++)
42+
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve))
43+
return {
44+
origin: `http://127.0.0.1:${(server.address() as AddressInfo).port}`,
45+
connections: () => connections,
46+
}
47+
}
48+
49+
describe('secureFetchWithPinnedIP connection reuse', () => {
50+
it('reuses one pinned connection across requests that share a pool', async () => {
51+
const server = await startServer((_req, res) => res.end('ok'))
52+
const pool = createPinnedConnectionPool()
53+
try {
54+
for (let index = 0; index < 3; index++) {
55+
const response = await secureFetchWithPinnedIP(server.origin, '127.0.0.1', {
56+
profile: 'configuredEndpoint',
57+
connectionPool: pool,
58+
})
59+
await expect(response.text()).resolves.toBe('ok')
60+
}
61+
expect(server.connections()).toBe(1)
62+
} finally {
63+
pool.destroy()
64+
}
65+
})
66+
67+
it('opens a fresh connection per request without a pool', async () => {
68+
const server = await startServer((_req, res) => res.end('ok'))
69+
for (let index = 0; index < 2; index++) {
70+
const response = await secureFetchWithPinnedIP(server.origin, '127.0.0.1', {
71+
profile: 'configuredEndpoint',
72+
})
73+
await response.text()
74+
}
75+
expect(server.connections()).toBe(2)
76+
})
77+
78+
it('never shares an agent between different pinned addresses', () => {
79+
const pool = createPinnedConnectionPool()
80+
try {
81+
const first = pool.agent(true, 'api.example.com', 443, '203.0.113.1')
82+
expect(pool.agent(true, 'api.example.com', 443, '203.0.113.1')).toBe(first)
83+
expect(pool.agent(true, 'api.example.com', 443, '203.0.113.2')).not.toBe(first)
84+
expect(pool.agent(true, 'other.example.com', 443, '203.0.113.1')).not.toBe(first)
85+
} finally {
86+
pool.destroy()
87+
}
88+
})
89+
})
90+
91+
describe('secureFetchWithPinnedIP compressed responses', () => {
92+
it('asks for compression only when requested and returns the decoded body', async () => {
93+
const encodings: (string | undefined)[] = []
94+
const server = await startServer((req, res) => {
95+
encodings.push(req.headers['accept-encoding'])
96+
if (req.headers['accept-encoding']?.includes('gzip')) {
97+
res.writeHead(200, { 'Content-Encoding': 'gzip' })
98+
res.end(gzipSync(Buffer.from('{"ok":true}')))
99+
} else {
100+
res.end('{"ok":true}')
101+
}
102+
})
103+
const compressed = await secureFetchWithPinnedIP(server.origin, '127.0.0.1', {
104+
profile: 'configuredEndpoint',
105+
acceptCompressed: true,
106+
})
107+
await expect(compressed.json()).resolves.toEqual({ ok: true })
108+
const plain = await secureFetchWithPinnedIP(server.origin, '127.0.0.1', {
109+
profile: 'configuredEndpoint',
110+
})
111+
await expect(plain.json()).resolves.toEqual({ ok: true })
112+
expect(encodings).toEqual(['gzip, deflate, br', undefined])
113+
})
114+
})
115+
116+
describe('secureFetchWithPinnedIP address pinning', () => {
117+
it.each([
118+
['a fresh agent', undefined],
119+
['a pooled agent', createPinnedConnectionPool()],
120+
])('pins the address on the request itself with %s', async (_label, pool) => {
121+
const server = await startServer((_req, res) => res.end('ok'))
122+
const request = vi.spyOn(http, 'request')
123+
try {
124+
const response = await secureFetchWithPinnedIP(
125+
server.origin.replace('127.0.0.1', 'localhost'),
126+
'127.0.0.1',
127+
{
128+
profile: 'configuredEndpoint',
129+
connectionPool: pool,
130+
}
131+
)
132+
await response.text()
133+
const options = request.mock.calls[0]?.[0] as http.RequestOptions
134+
const lookup = options.lookup as unknown as (
135+
hostname: string,
136+
options: object,
137+
callback: (error: Error | null, address: string, family: number) => void
138+
) => void
139+
const resolved = await new Promise<string>((resolve) =>
140+
lookup('localhost', {}, (_error, address) => resolve(address))
141+
)
142+
expect(resolved).toBe('127.0.0.1')
143+
} finally {
144+
request.mockRestore()
145+
pool?.destroy()
146+
}
147+
})
148+
})

‎apps/sim/lib/knowledge/application/chat.test.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ describe('organization Search Assistant chat', () => {
140140
userId: 'member-1',
141141
organizationId: 'org-1',
142142
chatId: 'private-chat',
143-
goRoute: '/api/mothership/execute',
143+
goRoute: '/api/mothership',
144144
interactive: false,
145145
autoExecuteTools: true,
146146
secretActorUserId: null,

‎apps/sim/lib/knowledge/application/chat.ts‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -172,7 +172,7 @@ export const organizationSearchChat: OperationUseCase<
172172
organizationId,
173173
chatId,
174174
simRequestId: messageId,
175-
goRoute: '/api/mothership/execute',
175+
goRoute: '/api/mothership',
176176
interactive: false,
177177
autoExecuteTools: true,
178178
secretActorUserId: null,

‎apps/sim/lib/knowledge/search/diagnostics.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,13 @@ export type SearchStage =
6464
| 'source_overview.searchable'
6565
| 'access_batch.connectors'
6666
| 'access_batch.live_proof'
67+
| 'live.policies'
68+
| 'live.accounts'
69+
| 'live.resolve'
70+
| 'live.session'
71+
| 'live.search'
72+
| 'live.verify'
73+
| 'live.read'
6774

6875
/** Fixed, content-free fields. Never pass queries, filters, document identities, SQL, or errors. */
6976
export interface SearchDiagnosticMetadata {

‎apps/sim/lib/sim-search/live/README.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -48,7 +48,7 @@ Selected labels are alternatives. Source date/query settings are authoritative r
4848

4949
### Google Calendar
5050

51-
Member mode searches calendars accessible through the member's account. Service mode verifies the account's primary-calendar identity, Directory customer, and source user selection before delegating to that same Workspace user. Selected calendar IDs constrain retrieval and source verification; `primary` means that member's primary calendar. The admin picker browses the delegated administrator's calendar list and stores `primary` as a per-member alias. Service search delegates with only `calendar.events.readonly` (plus `admin.directory.user.readonly` for the Directory check), the same scope as the indexed crawl; all-day events take the calendar's time zone from the events list response. The picker additionally needs `calendar.readonly` ([CalendarList authorization](https://developers.google.com/workspace/calendar/api/v3/reference/calendarList/list)).
51+
Member mode searches calendars accessible through the member's account. A meeting visible in several calendars is returned once, preferring the member's primary calendar, and date-bounded agendas are ordered by scheduled start across calendars. Service mode verifies the account's primary-calendar identity, Directory customer, and source user selection before delegating to that same Workspace user. Selected calendar IDs constrain retrieval and source verification; `primary` means that member's primary calendar. The admin picker browses the delegated administrator's calendar list and stores `primary` as a per-member alias. Service search delegates with only `calendar.events.readonly` (plus `admin.directory.user.readonly` for the Directory check), the same scope as the indexed crawl; all-day events take the calendar's time zone from the events list response. The picker additionally needs `calendar.readonly` ([CalendarList authorization](https://developers.google.com/workspace/calendar/api/v3/reference/calendarList/list)).
5252

5353
Each event must exist under the source's delegated token, be in an allowed calendar, not be cancelled, and overlap the source's configured rolling time window. The existing default is 30 days before and after the request. A stable UTC-day envelope around that window is intersected with the user's date bounds in the provider query, ensuring recurring events expand and query bounds stay stable between pages. The exact rolling source window is still checked for each result. Nonoverlapping date ranges return no results; continuations spanning a UTC-day change may require a fresh search. A source search query is checked with Calendar's event search and exact event-ID matching. All-day events use the calendar's timezone. Attendee details follow the source's include-attendees setting. The member's API access still determines which event details they can see.
5454

@@ -58,7 +58,7 @@ Member mode only. Search uses the connected user's Slack real-time search grant
5858

5959
### GitHub
6060

61-
Member mode searches issues, code, and repositories permitted by the connected token. Explicit repository/organization/user qualifiers narrow the user's query. Default discovery is bounded to up to 100 affiliated repositories, and provider pagination/search caps still apply.
61+
Member mode searches issues, code, and repositories permitted by the connected token. Explicit repository/organization/user qualifiers narrow the user's query. Default discovery is bounded to up to 100 affiliated repositories, sent as `repo:` qualifiers in at most four batches per search kind. Code batches stay under code search's 1,000-byte query limit, issue batches are larger, and a long query that cannot fit every repository reports the searched subset. Date bounds become one `updated:start..end` range, since GitHub ORs repeated qualifiers. Provider pagination/search caps still apply. REST code search returns no file dates and accepts no date qualifier, so date-filtered searches cover issues and pull requests only.
6262

6363
In service mode, an administrator connects a GitHub App installation and selects repositories one by one in Sources. Each source pins the provider-verified repository ID and may narrow code files by directory and extension. Search queries the member's own GitHub connection with `repo:` qualifiers drawn only from active sources. For each candidate, Sim checks that the current App installation still covers that repository, mints a repository-scoped read token, and compares repository and owner IDs returned under both the App and member tokens. It then checks the per-repository code filters. Reads use the member token and repeat these checks. A personal repository outside the selected sources is never searched, even if the member can access it. GitHub REST code search covers the default branch; live Sources therefore do not offer a branch setting.
6464

0 commit comments

Comments
 (0)