Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions apps/sim/lib/knowledge/mcp/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import {
import { liveCitationId } from '@/lib/knowledge/search/citation'
import { toolError } from '@/lib/mcp/tool-result'
import { readLiveDocument, searchLiveKnowledge } from '@/lib/sim-search/live/application'
import { LiveReadError } from '@/lib/sim-search/live/read-error'
import { v2CaughtOrchestrationError } from '@/app/api/v2/lib/response'
import { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'

Expand Down Expand Up @@ -67,6 +68,7 @@ export function createKnowledgeMcpServer(context: KnowledgeMcpContext): McpServe
return result
} catch (error) {
if (signal.aborted) outcome = 'cancelled'
if (error instanceof LiveReadError) return toolError(error.message)
const response = v2CaughtOrchestrationError(error)
if (response) {
const body: unknown = await response.json()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import {
import type { BaseServerTool, ServerToolContext } from '@/lib/mothership/tools/server/base-tool'
import { connectorDisplayName } from '@/lib/sim-search/connectors'
import { readLiveDocument, searchLiveKnowledge } from '@/lib/sim-search/live/application'
import { LiveReadError } from '@/lib/sim-search/live/read-error'
import { projectResolvedSecretModelContent } from '@/executor/utils/resolved-secret-content-projection'

const logger = createLogger('WorkspaceSearchTool')
Expand Down Expand Up @@ -165,8 +166,16 @@ export const readDocumentServerTool: BaseServerTool = {
}
} catch (error) {
logger.error('Document read failed', { error })
if (error instanceof LiveReadError)
return {
success: false,
retryable: error.retryable,
...(error.retryAfterSeconds ? { retryAfterSeconds: error.retryAfterSeconds } : {}),
message: error.message,
}
return {
success: false,
retryable: false,
message:
error instanceof z.ZodError
? 'Invalid document arguments'
Expand Down
57 changes: 57 additions & 0 deletions apps/sim/lib/sim-search/live/application.test.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { ErrorCode, McpError } from '@modelcontextprotocol/sdk/types.js'
import { dbChainMockFns, resetDbChainMock } from '@sim/testing'
import { createSessionPrincipal } from '@sim/testing/factories/principal.factory'
import { setEnv } from '@sim/testing/mocks/env.mock'
Expand Down Expand Up @@ -769,6 +770,62 @@ describe('authorized live retrieval', () => {
).rejects.toThrow('Revoked')
expect(mocks.read).not.toHaveBeenCalled()
})
it.each([
[
'a missing or non-readable document',
new NativeSearchError('unavailable', 'Provider request failed (404).', undefined, 404),
{ code: 'not_found', message: expect.stringContaining('Search again') },
],
[
'a revoked grant',
new NativeSearchError('reconnect', 'The provider denied access.'),
{ code: 'unauthorized', message: expect.stringContaining('Reconnect Google Drive') },
],
[
'a provider rate limit',
new NativeSearchError('rate_limited', 'Provider rate limit reached.', 30),
{ retryable: true, retryAfterSeconds: 30 },
],
[
'a provider outage',
new NativeSearchError('unavailable', 'Provider request failed (503).', undefined, 503),
{ retryable: true, message: expect.stringContaining('Try again') },
],
[
'an MCP request timeout',
new McpError(ErrorCode.RequestTimeout, 'TimeoutError'),
{ retryable: true, message: expect.stringContaining('took too long') },
],
[
'a provider request timeout response',
new NativeSearchError('unavailable', 'Provider request failed (408).', undefined, 408),
{ retryable: true, message: expect.stringContaining('timed out') },
],
[
'a native request socket timeout',
Object.assign(new Error('Request timed out after 10000ms'), { code: 'ETIMEDOUT' }),
{ retryable: true, message: expect.stringContaining('Try again') },
],
[
'a native request dispatcher timeout',
Object.assign(new Error('Headers Timeout Error'), { code: 'UND_ERR_HEADERS_TIMEOUT' }),
{ retryable: true, message: expect.stringContaining('Try again') },
],
])('classifies %s during a read so the caller can act on it', async (_, failure, expected) => {
const search = await searchLiveKnowledge.execute({ principal, input })
mocks.read.mockRejectedValueOnce(failure)
await expect(
readLiveDocument.execute({
principal,
input: {
workspaceId: 'workspace',
documentId: search.results[0]?.documentId ?? '',
limit: 1,
resultSecretRegistry: new ResolvedSecretTraceRegistry([]),
},
})
).rejects.toMatchObject(expected)
})
it.each([
['a stop reason', 'user_stop:test'],
['an AbortError', new DOMException('The operation was aborted.', 'AbortError')],
Expand Down
72 changes: 68 additions & 4 deletions apps/sim/lib/sim-search/live/application.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { ErrorCode, McpError } from '@modelcontextprotocol/sdk/types.js'
import { requirePrincipalSubjectUserId } from '@sim/auth/principal'
import { safeCompare } from '@sim/security/compare'
import { hmacSha256Hex } from '@sim/security/hmac'
Expand All @@ -17,6 +18,7 @@ import {
} from '@/lib/api/contracts/mothership-assistant-tools'
import { canonicalJson, fingerprint, instantScopePart } from '@/lib/api/cursor-binding'
import { env } from '@/lib/core/config/env'
import { isRetryableNetworkError } from '@/lib/core/errors/retryable-infrastructure'
import { OrchestrationError } from '@/lib/core/orchestration/types'
import {
type ResourceOwner,
Expand All @@ -25,6 +27,7 @@ import {
} from '@/lib/core/resource-scope'
import { createPinnedConnectionPool } from '@/lib/core/security/input-validation.server'
import { mapWithConcurrency } from '@/lib/core/utils/concurrency'
import { MANAGED_MCP_CONNECTORS } from '@/lib/credential-groups/managed-mcp-connectors'
import { requireOrganizationSearchAvailable } from '@/lib/knowledge/access/availability'
import { defineAuthorizedKnowledgeUseCase } from '@/lib/knowledge/application/authorized-knowledge-use-case'
import { resolveKnowledgeOwnerContext } from '@/lib/knowledge/application/contexts'
Expand All @@ -33,6 +36,7 @@ import { measureSearchStage } from '@/lib/knowledge/search/diagnostics'
import { RRF_K } from '@/lib/knowledge/search/recency'
import { matchPassage } from '@/lib/knowledge/search/snippet'
import { isKnowledgeSourceUrl } from '@/lib/knowledge/search/source-url'
import { connectorDisplayName } from '@/lib/sim-search/connectors'
import {
type LiveAccountSession,
openLiveAccountSession,
Expand All @@ -51,10 +55,12 @@ import {
withImpliedListingBound,
} from '@/lib/sim-search/live/dates'
import { NativeSearchError } from '@/lib/sim-search/live/http'
import { isManagedSearchMcpProvider } from '@/lib/sim-search/live/managed-mcp-config'
import { joinMessages } from '@/lib/sim-search/live/pages'
import { loadLiveSearchPolicies } from '@/lib/sim-search/live/policy-store'
import { LIVE_SEARCH_PROVIDER_IDS } from '@/lib/sim-search/live/provider-catalog'
import { liveSearchGuidance } from '@/lib/sim-search/live/providers'
import { LiveReadError } from '@/lib/sim-search/live/read-error'
import type { LiveAccount, NativeDocument } from '@/lib/sim-search/live/types'
import { projectResolvedSecretModelContent } from '@/executor/utils/resolved-secret-content-projection'
import type { ResolvedSecretTraceRegistry } from '@/executor/utils/resolved-secret-trace-registry'
Expand Down Expand Up @@ -746,6 +752,58 @@ export type LiveReadInput = ResourceOwner & {
signal?: AbortSignal
}

/** Milliseconds one document read may spend on provider calls, after account resolution. */
const READ_DEADLINE_MS = 15_000

function liveProviderName(provider: string): string {
return isManagedSearchMcpProvider(provider)
? MANAGED_MCP_CONNECTORS[provider].name
: connectorDisplayName(provider)
}

/**
* Classifies a provider failure during a read so the caller can act on it: a missing or
* non-readable document, a grant to reconnect, or a transient failure worth retrying.
* Classified and unrecognized errors pass through unchanged.
*/
function liveReadFailure(error: unknown, provider: string, deadline?: AbortSignal): unknown {
if (error instanceof OrchestrationError) return error
const name = liveProviderName(provider)
if (error instanceof NativeSearchError) {
if (error.status === 'reconnect')
return new OrchestrationError(
'unauthorized',
`Reconnect ${name} to read this document. ${error.message}`
)
if (error.status === 'rate_limited')
return new LiveReadError(error.message, true, error.retryAfterSeconds)
if (error.status === 'timeout' || error.httpStatus === 408)
return new LiveReadError(`${name} timed out reading this document. Try again.`, true)
if (error.httpStatus === 404 || error.httpStatus === 410)
return new OrchestrationError(
'not_found',
`${name} could not find this document: it was deleted, moved, or is not a readable page. Search again or read a different result.`
)
if (error.httpStatus !== undefined && error.httpStatus >= 500)
Comment thread
waleedlatif1 marked this conversation as resolved.
return new LiveReadError(
`${name} is temporarily unavailable (${error.httpStatus}). Try again shortly.`,
true
)
return new LiveReadError(error.message, false)
}
if (deadline?.aborted || (error instanceof McpError && error.code === ErrorCode.RequestTimeout))
return new LiveReadError(
`${name} took too long to return this document. Try again, or read a different result.`,
true
)
if (isRetryableNetworkError(error))
return new LiveReadError(
`${name} did not respond while reading this document. Try again shortly.`,
true
)
return error
Comment thread
waleedlatif1 marked this conversation as resolved.
Comment thread
waleedlatif1 marked this conversation as resolved.
}

export const readLiveDocument = defineAuthorizedKnowledgeUseCase({
operation: knowledgeOperations.readDocument,
resolveContext: ({ input }: { input: LiveReadInput }) => resolveKnowledgeOwnerContext(input),
Expand All @@ -765,15 +823,19 @@ export const readLiveDocument = defineAuthorizedKnowledgeUseCase({
(input.filters?.documentIds && !input.filters.documentIds.includes(input.documentId))
)
throw new OrchestrationError('not_found', 'Document is outside the selected search filters')
/** Cancellation by the caller passes through unclassified. */
const readFailure = (error: unknown, deadline?: AbortSignal) =>
input.signal?.aborted ? error : liveReadFailure(error, reference.provider, deadline)
const [resolved, policies] = await Promise.all([
resolveLiveAccount(input, userId, reference.account),
loadLiveSearchPolicies(input),
])
]).catch((error: unknown) => {
throw readFailure(error)
})
if (resolved.account.provider !== reference.provider)
throw new OrchestrationError('not_found', 'Document account changed')
const signal = input.signal
? AbortSignal.any([input.signal, AbortSignal.timeout(15_000)])
: AbortSignal.timeout(15_000)
const deadline = AbortSignal.timeout(READ_DEADLINE_MS)
const signal = input.signal ? AbortSignal.any([input.signal, deadline]) : deadline
const pool = createPinnedConnectionPool()
let document: NativeDocument
let session: LiveAccountSession | undefined
Expand Down Expand Up @@ -805,6 +867,8 @@ export const readLiveDocument = defineAuthorizedKnowledgeUseCase({
'not_found',
'Document is outside your organization’s search scope'
)
} catch (error) {
throw readFailure(error, deadline)
} finally {
try {
await session?.close()
Expand Down
26 changes: 26 additions & 0 deletions apps/sim/lib/sim-search/live/atlassian.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -155,4 +155,30 @@ describe('Confluence live documents', () => {
expect(page.documents[0]?.accessMetadata).toEqual({ spaceKey: 'ENG' })
expect(page.documents[2]?.accessMetadata).toEqual({ spaceKey: 'ENG' })
})

it('drops native CQL matches that the page and blog post reads cannot open', async () => {
const api = client({
'/ex/confluence/cloud/wiki/rest/api/search': {
results: [
{ content: { id: '123', type: 'page', title: 'Runbook' } },
{ content: { id: '77', type: 'attachment', title: 'diagram.png' } },
{ content: { id: '78', type: 'comment', title: 'Re: Runbook' } },
{ content: { id: '79', type: 'whiteboard', title: 'Planning' } },
{ content: { id: '80', type: 'folder', title: 'Archive' } },
],
_links: {},
},
})
const page = await searchAtlassian(api, 'confluence', {
query: '',
native: { provider: 'confluence', query: 'title ~ "runbook"' },
limit: 10,
scopes: [],
})
expect(page.documents.map(({ id, kind }) => ({ id, kind }))).toEqual([
{ id: '123', kind: 'page' },
])
expect(page.message).toContain('cannot be read')
expect(page.partial).toBe(true)
})
})
31 changes: 26 additions & 5 deletions apps/sim/lib/sim-search/live/atlassian.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,7 +65,16 @@ function issue(row: Record<string, unknown>, cloudId: string, site: string): Nat
author: string(object(fields.creator).displayName),
}
}
function page(row: Record<string, unknown>, cloudId: string, site: string): NativeDocument {
/**
* Maps a CQL search hit to a readable document, or `undefined` for kinds the v2 page and blog
* post reads cannot open. Native CQL can match attachments, comments, whiteboards, folders,
* databases, and users; returning those would hand the caller a reference whose read is a 404.
*/
function page(
row: Record<string, unknown>,
cloudId: string,
site: string
): NativeDocument | undefined {
if (string(row.entityType) === 'space') {
const space = object(row.space)
return {
Expand All @@ -83,9 +92,11 @@ function page(row: Record<string, unknown>, cloudId: string, site: string): Nati
const links = object(content._links)
const version = object(content.version)
const spaceKey = string(object(content.space).key)
const kind = string(content.type)
if ((kind !== 'page' && kind !== 'blogpost') || !string(content.id)) return undefined
return {
id: string(content.id),
kind: string(content.type) === 'blogpost' ? 'blogpost' : 'page',
kind,
...(spaceKey ? { accessMetadata: { spaceKey } } : {}),
container: cloudId,
title: string(content.title) || string(row.title),
Expand Down Expand Up @@ -164,6 +175,7 @@ export async function searchAtlassian(
)
return {
documents: array(data.issues).map((row) => issue(row, cloudId, origin)),
excluded: false,
next: string(data.nextPageToken) || undefined,
}
}
Expand All @@ -183,22 +195,31 @@ export async function searchAtlassian(
})
)
const next = string(object(data._links).next)
const rows = array(data.results)
const documents = rows
.map((row) => page(row, cloudId, origin))
.filter((document) => document !== undefined)
return {
documents: array(data.results).map((row) => page(row, cloudId, origin)),
documents,
excluded: documents.length < rows.length,
next: next
? (new URL(next, 'https://api.atlassian.com').searchParams.get('cursor') ?? undefined)
: undefined,
}
})
)
const documents = interleaveByRank(pages.map((result) => result.documents))
const excluded = pages.some((result) => result.excluded)
Comment thread
waleedlatif1 marked this conversation as resolved.
return {
documents,
partial: !input.native?.project && allSites.length > selected.length,
partial: excluded || (!input.native?.project && allSites.length > selected.length),
hasMore: pages.some((result) => Boolean(result.next)),
nextCursor: single ? pages[0]?.next : undefined,
message:
'Searches up to four accessible Atlassian sites. For a specific site and pagination, set project to its cloud ID.',
'Searches up to four accessible Atlassian sites. For a specific site and pagination, set project to its cloud ID.' +
(excluded
? ' Matches other than pages, blog posts, and spaces (attachments, comments, whiteboards, folders, databases) were excluded because they cannot be read.'
: ''),
}
}

Expand Down
4 changes: 3 additions & 1 deletion apps/sim/lib/sim-search/live/managed-mcp.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -739,7 +739,9 @@ describe('managed Search operation sessions', () => {
.where(eq(credential.id, actors[0].credentialId))
}
try {
await expect(read(found.results[0].documentId)).rejects.toMatchObject({ status: 'reconnect' })
await expect(read(found.results[0].documentId)).rejects.toMatchObject({
code: 'unauthorized',
})
expect(openTransports.size).toBe(0)
} finally {
onTool = undefined
Expand Down
15 changes: 15 additions & 0 deletions apps/sim/lib/sim-search/live/read-error.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
/**
* A provider-side read failure whose message tells the caller what to do next, and whether the
* same read may succeed if retried. Messages are written by Search code, never copied from a
* provider response body.
*/
export class LiveReadError extends Error {
constructor(
message: string,
readonly retryable: boolean,
readonly retryAfterSeconds?: number
) {
super(message)
this.name = 'LiveReadError'
}
}
Loading