Skip to content

Commit 9cfa049

Browse files
authored
fix(search): complete sync failure diagnostic context (#7898)
1 parent 90a05d9 commit 9cfa049

7 files changed

Lines changed: 207 additions & 10 deletions

File tree

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
/** @vitest-environment node */
2+
import { DrizzleQueryError } from 'drizzle-orm/errors'
3+
import { beforeEach, describe, expect, it, vi } from 'vitest'
4+
5+
const { refresh } = vi.hoisted(() => ({ refresh: vi.fn() }))
6+
vi.mock('@/lib/knowledge/connectors/external-group-sync', () => ({
7+
refreshConnectorDirectory: refresh,
8+
}))
9+
10+
import { executeDirectorySyncJob } from '@/background/knowledge-connector-directory-sync'
11+
import { GoogleDriveApiError } from '@/connectors/google-drive/google-drive-errors'
12+
import { ConnectorDirectoryError } from '@/connectors/source-error'
13+
14+
const PAYLOAD = { connectorId: 'connector-1', requestId: 'request-1' }
15+
16+
describe('directory sync worker diagnostics', () => {
17+
beforeEach(() => vi.clearAllMocks())
18+
19+
it('preserves successful directory outcomes', async () => {
20+
refresh.mockResolvedValueOnce('refreshed')
21+
await expect(executeDirectorySyncJob(PAYLOAD)).resolves.toEqual({ outcome: 'refreshed' })
22+
expect(refresh).toHaveBeenCalledWith('connector-1', 'request-1')
23+
})
24+
25+
it('includes the Google operation and reason in the final task failure', async () => {
26+
refresh.mockRejectedValueOnce(
27+
new ConnectorDirectoryError('private group detail', {
28+
cause: new GoogleDriveApiError(403, ['forbidden'], 'directory.members.list'),
29+
})
30+
)
31+
const error = await executeDirectorySyncJob(PAYLOAD).catch((error: unknown) => error)
32+
expect(error).toMatchObject({
33+
message:
34+
'Directory permission sync failed (HTTP 403). Operation: directory.members.list. Google reason: forbidden. Group membership could not be fully verified.',
35+
})
36+
expect(error).not.toHaveProperty('cause')
37+
expect(String(error)).not.toContain('private')
38+
})
39+
40+
it('preserves database codes without exposing the raw driver cause to Trigger', async () => {
41+
refresh.mockRejectedValueOnce(
42+
new DrizzleQueryError(
43+
'select private SQL',
44+
['private value'],
45+
Object.assign(new Error('private driver detail'), { code: '57014' })
46+
)
47+
)
48+
const error = await executeDirectorySyncJob(PAYLOAD).catch((error: unknown) => error)
49+
expect(error).toMatchObject({ message: 'Database request failed (SQLSTATE 57014).' })
50+
expect(error).not.toHaveProperty('cause')
51+
expect(String(error)).not.toContain('private')
52+
})
53+
54+
it('preserves unclassified failures', async () => {
55+
const error = new Error('unexpected failure')
56+
refresh.mockRejectedValueOnce(error)
57+
await expect(executeDirectorySyncJob(PAYLOAD)).rejects.toBe(error)
58+
})
59+
60+
it('retains the database code when a directory refresh wraps the driver failure', async () => {
61+
refresh.mockRejectedValueOnce(
62+
new ConnectorDirectoryError('private wrapper', {
63+
cause: new DrizzleQueryError(
64+
'select private SQL',
65+
['private value'],
66+
Object.assign(new Error('private driver detail'), { code: '57014' })
67+
),
68+
})
69+
)
70+
const error = await executeDirectorySyncJob(PAYLOAD).catch((error: unknown) => error)
71+
expect(error).toMatchObject({
72+
message:
73+
'Directory permission sync failed. Error code: 57014. Group membership could not be fully verified.',
74+
})
75+
expect(error).not.toHaveProperty('cause')
76+
expect(String(error)).not.toContain('private')
77+
})
78+
})

‎apps/sim/background/knowledge-connector-directory-sync.ts‎

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import { createLogger } from '@sim/logger'
22
import { task } from '@trigger.dev/sdk'
3+
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
34
import {
45
assertDirectorySyncPayload,
56
DIRECTORY_SYNC_CONCURRENCY,
@@ -14,9 +15,16 @@ const logger = createLogger('TriggerKnowledgeConnectorDirectorySync')
1415
export async function executeDirectorySyncJob(payload: unknown) {
1516
const { connectorId, requestId } = assertDirectorySyncPayload(payload)
1617
logger.info(`[${requestId}] Starting directory refresh: ${connectorId}`)
17-
const outcome = await refreshConnectorDirectory(connectorId, requestId)
18-
logger.info(`[${requestId}] Directory refresh finished`, { connectorId, outcome })
19-
return { outcome }
18+
try {
19+
const outcome = await refreshConnectorDirectory(connectorId, requestId)
20+
logger.info(`[${requestId}] Directory refresh finished`, { connectorId, outcome })
21+
return { outcome }
22+
} catch (error) {
23+
const diagnostic = getConnectorFailureDiagnostic(error)
24+
if (!diagnostic) throw error
25+
logger.error(`[${requestId}] Directory refresh failed`, { connectorId, diagnostic })
26+
throw new Error(diagnostic.message)
27+
}
2028
}
2129

2230
export const knowledgeConnectorDirectorySync = task({

‎apps/sim/lib/knowledge/connectors/connector-error.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -125,13 +125,14 @@ export function getConnectorFailureDiagnostic(error: unknown): ConnectorFailureD
125125
const context = sourceError?.diagnostic
126126
if (directoryError) {
127127
const status = diagnostic?.status ? ` (HTTP ${diagnostic.status})` : ''
128+
const code = diagnostic?.code ? ` Error code: ${diagnostic.code}.` : ''
128129
const reason = context?.reasons.length ? ` Google reason: ${context.reasons.join(', ')}.` : ''
129130
return {
130131
...diagnostic,
131132
...context,
132133
category: diagnostic?.category ?? 'directory',
133134
phase: 'directory',
134-
message: `Directory permission sync failed${status}.${context ? ` Operation: ${context.operation}.` : ''}${reason} Group membership could not be fully verified.`,
135+
message: `Directory permission sync failed${status}.${context ? ` Operation: ${context.operation}.` : ''}${reason}${code} Group membership could not be fully verified.`,
135136
}
136137
}
137138
if (!diagnostic || !context) return diagnostic

‎apps/sim/lib/knowledge/connectors/external-group-sync.test.ts‎

Lines changed: 28 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@
44
import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing'
55
import { beforeEach, describe, expect, it, vi } from 'vitest'
66
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
7+
import { GoogleDriveApiError } from '@/connectors/google-drive/google-drive-errors'
78
import type { ConnectorDirectory } from '@/connectors/types'
89

910
const { mockResolveTokenUserId, mockResolveToken, mockOpenDirectory, mockAvailability } =
@@ -254,7 +255,8 @@ describe('refreshConnectorDirectory', () => {
254255
)
255256
expect(dbChainMockFns.set).toHaveBeenCalledWith(
256257
expect.objectContaining({
257-
lastSyncError: 'Directory refresh failed: 403',
258+
lastSyncError:
259+
'Directory refresh failed: Directory permission sync failed. Group membership could not be fully verified.',
258260
})
259261
)
260262
expect(dbChainMockFns.delete).not.toHaveBeenCalled()
@@ -274,6 +276,31 @@ describe('refreshConnectorDirectory', () => {
274276
expect(dbChainMockFns.set.mock.calls.some(([value]) => 'lastSyncedAt' in value)).toBe(false)
275277
})
276278

279+
it('persists the nested Google reason for scheduled directory failures', async () => {
280+
queueTableRows(schemaMock.knowledgeConnector, [connectorRow()])
281+
const providerError = new GoogleDriveApiError(403, ['forbidden'], 'directory.members.list')
282+
mockOpenDirectory.mockResolvedValue(
283+
directory({ listGroupMembers: vi.fn().mockRejectedValue(providerError) })
284+
)
285+
286+
const failure = await refreshConnectorDirectory('connector-1', 'req-1').catch(
287+
(error: unknown) => error
288+
)
289+
expect(getConnectorFailureDiagnostic(failure)).toMatchObject({
290+
status: 403,
291+
operation: 'directory.members.list',
292+
reasons: ['forbidden'],
293+
phase: 'directory',
294+
})
295+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
296+
expect.objectContaining({
297+
lastSyncError:
298+
'Directory refresh failed: Directory permission sync failed (HTTP 403). Operation: directory.members.list. Google reason: forbidden. Group membership could not be fully verified.',
299+
})
300+
)
301+
expect(dbChainMockFns.set.mock.calls.some(([value]) => 'lastSyncedAt' in value)).toBe(false)
302+
})
303+
277304
it('clears a previous directory error after a successful refresh', async () => {
278305
queueTableRows(schemaMock.knowledgeConnector, [
279306
connectorRow({ lastSyncError: 'Directory refresh failed: 403' }),

‎apps/sim/lib/knowledge/connectors/external-group-sync.ts‎

Lines changed: 11 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ import {
2828
resolveConnectorTokenUserId,
2929
syncContextForToken,
3030
} from '@/lib/knowledge/connectors/access-token'
31+
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
3132
import { RUNNABLE_CONNECTOR_STATUSES } from '@/lib/knowledge/connectors/sync-lock'
3233
import { isRateLimitError } from '@/lib/knowledge/documents/utils'
3334
import { CONNECTOR_REGISTRY } from '@/connectors/registry.server'
@@ -189,11 +190,13 @@ export async function syncExternalDirectoryGroups(input: {
189190
if (isRateLimitError(error)) throw error
190191
keptStale += 1
191192
firstError ??= toError(error)
193+
const diagnostic = getConnectorFailureDiagnostic(error)
192194
logger.warn('Keeping last-known-good membership for a group that failed to enumerate', {
193195
workspaceId,
194196
providerId,
195197
externalGroupId: group.id,
196-
error: getErrorMessage(error),
198+
error: diagnostic?.message ?? getErrorMessage(error),
199+
diagnostic,
197200
})
198201
continue
199202
}
@@ -378,10 +381,12 @@ export async function refreshMirroredDirectory(input: {
378381
})
379382
return result.skipped ? 'skipped' : 'refreshed'
380383
} catch (error) {
384+
const diagnostic = getConnectorFailureDiagnostic(error)
381385
logger.error('Directory refresh failed; serving last-known-good group membership', {
382386
workspaceId,
383387
connector: connectorConfig.id,
384-
error: getErrorMessage(error),
388+
error: diagnostic?.message ?? getErrorMessage(error),
389+
diagnostic,
385390
})
386391
throw new ConnectorDirectoryError(`${DIRECTORY_ERROR_PREFIX}${getErrorMessage(error)}`, {
387392
cause: error,
@@ -506,7 +511,10 @@ export async function refreshConnectorDirectory(
506511
}
507512
return outcome
508513
} catch (error) {
509-
await recordError(getErrorMessage(error))
514+
const diagnostic = getConnectorFailureDiagnostic(error)
515+
await recordError(
516+
diagnostic ? `${DIRECTORY_ERROR_PREFIX}${diagnostic.message}` : getErrorMessage(error)
517+
)
510518
throw error
511519
}
512520
}

‎apps/sim/lib/knowledge/documents/document-processing-source.test.ts‎

Lines changed: 59 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -746,6 +746,65 @@ describe('processDocumentAsync write guards', () => {
746746
expect(guardForStatusWrite('failed')).toBeDefined()
747747
})
748748

749+
it('records the failed embedding batch without exposing SQL, content or vectors', async () => {
750+
armProviderSource()
751+
dbChainMockFns.limit.mockResolvedValueOnce([{ id: 'document-1' }])
752+
mockProcessDocument.mockResolvedValueOnce({
753+
chunks: [{ text: 'private-content', metadata: { startIndex: 0, endIndex: 15 } }],
754+
metadata: { chunkCount: 1, tokenCount: 3, characterCount: 15 },
755+
})
756+
mockGenerateEmbeddings.mockResolvedValueOnce({
757+
embeddings: [[0.123456789]],
758+
billableTokens: 0,
759+
modelName: 'text-embedding-3-small',
760+
pricingId: 'text-embedding-3-small',
761+
})
762+
const databaseError = new DrizzleQueryError(
763+
'insert private SQL',
764+
['private-content', [0.123456789]],
765+
Object.assign(new Error('canceling statement due to statement timeout'), { code: '57014' })
766+
)
767+
dbChainMockFns.values.mockRejectedValueOnce(databaseError)
768+
769+
await expect(
770+
processDocumentAsync(
771+
'knowledge-base-1',
772+
'document-1',
773+
{
774+
filename: 'a.txt',
775+
fileUrl: 'https://example.com/a.txt',
776+
fileSize: 15,
777+
mimeType: 'text/plain',
778+
},
779+
{},
780+
BILLING_ATTRIBUTION
781+
)
782+
).rejects.toBe(databaseError)
783+
784+
expect(mockLogError).toHaveBeenCalledWith('[document-1] Failed to insert embedding batch', {
785+
knowledgeBaseId: 'knowledge-base-1',
786+
operation: 'embedding.insert',
787+
batchNumber: 1,
788+
batchSize: 1,
789+
totalChunks: 1,
790+
embeddingModel: 'text-embedding-3-small',
791+
embeddingDimensions: 1536,
792+
elapsedMs: expect.any(Number),
793+
diagnostic: {
794+
category: 'database',
795+
code: '57014',
796+
message: 'Database request failed (SQLSTATE 57014).',
797+
},
798+
})
799+
const logs = JSON.stringify(mockLogError.mock.calls)
800+
expect(logs).not.toContain('private')
801+
expect(logs).not.toContain('0.123456789')
802+
expect(guardForStatusWrite('failed')).toBeDefined()
803+
expect(
804+
dbChainMockFns.set.mock.calls.some(([value]) => value.processingStatus === 'completed')
805+
).toBe(false)
806+
})
807+
749808
it('accepts a legacy queuedAt-only payload only while the row has no token', async () => {
750809
dbChainMockFns.limit
751810
.mockResolvedValueOnce([PERSISTED_CONTEXT])

‎apps/sim/lib/knowledge/documents/service.ts‎

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1889,9 +1889,25 @@ export async function processDocumentAsync(
18891889
}
18901890

18911891
logger.info(`[${documentId}] Inserting ${embeddingRecords.length} embeddings`)
1892-
for (const batch of batches) {
1892+
for (const [batchIndex, batch] of batches.entries()) {
18931893
signal.throwIfAborted()
1894-
await tx.insert(embedding).values(batch)
1894+
const insertStartedAt = Date.now()
1895+
try {
1896+
await tx.insert(embedding).values(batch)
1897+
} catch (error) {
1898+
logger.error(`[${documentId}] Failed to insert embedding batch`, {
1899+
knowledgeBaseId,
1900+
operation: 'embedding.insert',
1901+
batchNumber: batchIndex + 1,
1902+
batchSize: batch.length,
1903+
totalChunks: embeddingRecords.length,
1904+
embeddingModel: kbEmbeddingModel,
1905+
embeddingDimensions: kbEmbedding.dimensions,
1906+
elapsedMs: Date.now() - insertStartedAt,
1907+
diagnostic: getConnectorFailureDiagnostic(error),
1908+
})
1909+
throw error
1910+
}
18951911
}
18961912
const provenanceRecords = embeddingRecords.flatMap((record, index) => {
18971913
const provenance = chunkProvenances[index]

0 commit comments

Comments
 (0)