Skip to content

Commit 99d6a9f

Browse files
authored
fix(knowledge): preserve live document processing during recovery (#7993)
* fix(knowledge): preserve live document processing during recovery * fix(knowledge): recover abandoned redelivery and cancel liveness requests * fix(knowledge): fence recovery against concurrent protection
1 parent f359031 commit 99d6a9f

17 files changed

Lines changed: 1445 additions & 503 deletions

‎apps/sim/lib/internal/mistral/client.test.ts‎

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,11 @@
33
*/
44
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
55

6+
const { errorLog } = vi.hoisted(() => ({ errorLog: vi.fn() }))
7+
vi.mock('@sim/logger', () => ({
8+
createLogger: () => ({ error: errorLog, info: vi.fn(), warn: vi.fn(), debug: vi.fn() }),
9+
}))
10+
611
const { fetchPinned, admit, settle, validate } = vi.hoisted(() => ({
712
fetchPinned: vi.fn(),
813
admit: vi.fn(),
@@ -63,9 +68,41 @@ describe('Mistral provider transport', () => {
6368
expect(settle).toHaveBeenCalledWith('success', undefined)
6469
})
6570

71+
it('records cooldown without waiting for a stalled 429 error body', async () => {
72+
const cancel = vi.fn()
73+
fetchPinned.mockResolvedValue(
74+
new Response(new ReadableStream({ cancel }), {
75+
status: 429,
76+
headers: { 'retry-after': '60' },
77+
})
78+
)
79+
await expect(submitMistralOcr('private-key', {})).rejects.toMatchObject({
80+
reason: 'rate_limit',
81+
retryAfterMs: 60_000,
82+
})
83+
expect(settle).toHaveBeenCalledWith('rate_limit', 60_000)
84+
expect(cancel).toHaveBeenCalledOnce()
85+
expect(fetchPinned).toHaveBeenCalledOnce()
86+
})
87+
88+
it('preserves provider rejection when its diagnostic body exceeds the byte limit', async () => {
89+
fetchPinned.mockResolvedValue(new Response('x'.repeat(70_000), { status: 400 }))
90+
await expect(submitMistralOcr('private-key', {})).rejects.toMatchObject({
91+
status: 400,
92+
body: { success: false, error: 'Mistral API error: HTTP 400' },
93+
})
94+
expect(errorLog).toHaveBeenCalledWith(
95+
'Mistral API error',
96+
expect.objectContaining({ status: 400, bodyFormat: 'unavailable' })
97+
)
98+
})
99+
66100
it('identifies provider request rejection without retaining echoed document contents', async () => {
67101
fetchPinned.mockResolvedValue(
68-
Response.json({ message: 'Sensitive fixture document text' }, { status: 400 })
102+
Response.json(
103+
{ type: 'invalid_request_error', code: 400, message: 'Sensitive fixture document text' },
104+
{ status: 400, headers: { 'x-request-id': 'ocr-request-123' } }
105+
)
69106
)
70107
await expect(submitMistralOcr('key', {})).rejects.toMatchObject({
71108
source: 'provider',
@@ -74,6 +111,16 @@ describe('Mistral provider transport', () => {
74111
})
75112
expect(fetchPinned).toHaveBeenCalledOnce()
76113
expect(settle).toHaveBeenCalledWith('failure', undefined)
114+
expect(errorLog).toHaveBeenCalledWith(
115+
'Mistral API error',
116+
expect.objectContaining({
117+
status: 400,
118+
providerRequestId: 'ocr-request-123',
119+
providerErrorCode: '400',
120+
providerErrorType: 'invalid_request_error',
121+
})
122+
)
123+
expect(JSON.stringify(errorLog.mock.calls)).not.toContain('Sensitive fixture document text')
77124
})
78125

79126
it.each([

‎apps/sim/lib/internal/mistral/client.ts‎

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -11,9 +11,10 @@ import {
1111
validateUrlWithDNS,
1212
} from '@/lib/core/security/input-validation.server'
1313
import { getMistralCapacityConfig, getMistralCapacityScope } from '@/lib/internal/mistral/capacity'
14+
import { getOcrResponseDiagnostic } from '@/lib/internal/mistral/error-diagnostics'
1415
import { MistralOperationError } from '@/lib/internal/mistral/errors'
1516
import { MISTRAL_OCR_REQUEST_POLICY } from '@/lib/knowledge/documents/ocr-request-policy'
16-
import { readBoundedHttpErrorBody, resolveRetryDelayMs } from '@/lib/knowledge/documents/utils'
17+
import { readBoundedHttpErrorPayload, resolveRetryDelayMs } from '@/lib/knowledge/documents/utils'
1718

1819
const logger = createLogger('MistralClient')
1920
const MISTRAL_ENDPOINT = 'https://api.mistral.ai/v1/ocr'
@@ -138,8 +139,13 @@ export async function submitMistralOcr(
138139
retryAfterMs,
139140
})
140141
}
141-
await readBoundedHttpErrorBody(response)
142-
logger.error('Mistral API error', { status: response.status })
142+
const payload = await readBoundedHttpErrorPayload(response)
143+
logger.error('Mistral API error', {
144+
provider: 'mistral',
145+
operation: 'ocr',
146+
status: response.status,
147+
...getOcrResponseDiagnostic(response.headers, payload.ok ? payload.body : ''),
148+
})
143149
throw new MistralOperationError(
144150
response.status,
145151
{ success: false, error: `Mistral API error: HTTP ${response.status}` },
Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,49 @@
1+
/** @vitest-environment node */
2+
import { describe, expect, it } from 'vitest'
3+
import { getOcrResponseDiagnostic } from '@/lib/internal/mistral/error-diagnostics'
4+
5+
describe('OCR error diagnostics', () => {
6+
it.each([
7+
{ code: 400, type: 'invalid_request_error', message: 'private document' },
8+
{ error: { code: 400, type: 'invalid_request_error', message: 'private document' } },
9+
])('projects only safe fields from a provider envelope', (body) => {
10+
expect(
11+
getOcrResponseDiagnostic(new Headers({ 'x-request-id': 'request_123' }), JSON.stringify(body))
12+
).toEqual({
13+
bodyFormat: 'json',
14+
providerRequestId: 'request_123',
15+
providerErrorCode: '400',
16+
providerErrorType: 'invalid_request_error',
17+
})
18+
})
19+
20+
it.each(['not json', '<html>private input</html>', '', 'null', '[]', '42'])(
21+
'tolerates an unavailable or non-object error: %s',
22+
(body) => {
23+
expect(getOcrResponseDiagnostic(new Headers(), body)).toMatchObject({
24+
providerRequestId: null,
25+
providerErrorCode: null,
26+
providerErrorType: null,
27+
})
28+
}
29+
)
30+
31+
it('does not trust arbitrary codes, error types, echoed input, or request IDs', () => {
32+
const diagnostic = getOcrResponseDiagnostic(
33+
new Headers({ 'x-request-id': 'secret '.repeat(30), 'apim-request-id': 'safe-123' }),
34+
JSON.stringify({
35+
code: 'private_document',
36+
type: { input: 'private_document' },
37+
message: 'private_document',
38+
param: 'secret-key',
39+
detail: { input: 'private_document' },
40+
})
41+
)
42+
expect(diagnostic).toMatchObject({
43+
providerRequestId: 'safe-123',
44+
providerErrorCode: 'unrecognized',
45+
providerErrorType: 'unrecognized',
46+
})
47+
expect(JSON.stringify(diagnostic)).not.toMatch(/private_document|secret-key/)
48+
})
49+
})
Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,57 @@
1+
const SAFE_ERROR_CODES = new Set([
2+
'invalid_request_error',
3+
'authentication_error',
4+
'permission_error',
5+
'rate_limit_error',
6+
'server_error',
7+
'unknown_model',
8+
'BadRequest',
9+
'InvalidRequest',
10+
'DeploymentNotFound',
11+
'ResourceNotFound',
12+
'OperationNotSupported',
13+
'Unauthorized',
14+
'Forbidden',
15+
'TooManyRequests',
16+
'InternalServerError',
17+
'ServiceUnavailable',
18+
])
19+
20+
function safeCode(value: unknown): string | null {
21+
if (value === undefined || value === null) return null
22+
if (typeof value === 'string' && SAFE_ERROR_CODES.has(value)) return value
23+
if (typeof value === 'number' && Number.isInteger(value) && value >= 400 && value <= 599)
24+
return String(value)
25+
return 'unrecognized'
26+
}
27+
28+
function safeRequestId(value: string | null): string | null {
29+
return value && /^[a-zA-Z0-9_-]{1,128}$/.test(value) ? value : null
30+
}
31+
32+
/** Projects bounded OCR error responses without retaining messages, document data or URLs. */
33+
export function getOcrResponseDiagnostic(headers: Pick<Headers, 'get'>, body: string) {
34+
const diagnostic = {
35+
providerRequestId:
36+
safeRequestId(headers.get('x-request-id')) ??
37+
safeRequestId(headers.get('apim-request-id')) ??
38+
safeRequestId(headers.get('x-ms-request-id')),
39+
bodyFormat: body ? 'non_json' : 'unavailable',
40+
providerErrorCode: null as string | null,
41+
providerErrorType: null as string | null,
42+
}
43+
if (!body) return diagnostic
44+
let parsed: unknown
45+
try {
46+
parsed = JSON.parse(body)
47+
} catch {
48+
return diagnostic
49+
}
50+
diagnostic.bodyFormat = 'json'
51+
if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) return diagnostic
52+
const error = 'error' in parsed ? parsed.error : parsed
53+
if (!error || typeof error !== 'object' || Array.isArray(error)) return diagnostic
54+
diagnostic.providerErrorCode = safeCode('code' in error ? error.code : undefined)
55+
diagnostic.providerErrorType = safeCode('type' in error ? error.type : undefined)
56+
return diagnostic
57+
}

0 commit comments

Comments
 (0)