Skip to content

Commit efd0535

Browse files
committed
fix(jobs): escalate database retries without spending the breaker, cover members mode
- connector and members-mode syncs climb the failure ladder by the failed-run streak read from their run logs; database failures never advance the auto-disable counter - classify 25P03 as capacity and ECONNREFUSED/EHOSTUNREACH/ENOTFOUND/EAI_AGAIN (with a query) as connection - in-process document processing keeps recording database failures as failed, documented - offer Retry for a pending document whose dispatch or deferred retry is past the retry API's grace - correct the processing task's retry-ceiling comment
1 parent 03ef2c9 commit efd0535

19 files changed

Lines changed: 659 additions & 54 deletions

File tree

‎apps/sim/app/workspace/[workspaceId]/knowledge/[id]/base.tsx‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,6 +49,7 @@ import {
4949
KNOWLEDGE_DOCUMENT_PROCESSING_STALE_THRESHOLD_MS,
5050
} from '@/lib/knowledge/constants'
5151
import {
52+
canRetryDocumentProcessing,
5253
type DocumentSortField,
5354
getDocumentIndexingStatus,
5455
type SortOrder,
@@ -1569,7 +1570,7 @@ export function KnowledgeBase({
15691570
}
15701571
onRetry={
15711572
contextMenuDocument &&
1572-
getDocumentIndexingStatus(contextMenuDocument) === 'failed' &&
1573+
canRetryDocumentProcessing(contextMenuDocument) &&
15731574
selectedDocumentCount === 1 &&
15741575
userPermissions.canEdit
15751576
? () => handleRetryDocument(contextMenuDocument.id)

‎apps/sim/background/knowledge-processing.ts‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -327,9 +327,10 @@ export const processDocument = task({
327327
machine: 'medium-2x',
328328
retry: {
329329
/**
330-
* The ceiling for database retries; `catchError` stops other failures at
331-
* `KB_CONFIG_MAX_ATTEMPTS`. An out-of-memory kill never reaches `catchError`,
332-
* so it may use the full ceiling.
330+
* The ceiling for thrown errors: database retries use all of it, and
331+
* `catchError` stops every other thrown error at `KB_CONFIG_MAX_ATTEMPTS`.
332+
* A crashed or timed-out run is not retried; an out-of-memory kill is
333+
* retried once, on the `outOfMemory` machine below.
333334
*/
334335
maxAttempts: backgroundRetryAttemptCeiling(DOCUMENT_PROCESSING_RETRY_POLICY),
335336
factor: envNumber(env.KB_CONFIG_RETRY_FACTOR, 2),

‎apps/sim/lib/api/contracts/knowledge/documents.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -293,6 +293,8 @@ export const documentDataSchema = z
293293
processingOutcome: z.literal('skipped').nullable().default(null),
294294
/** When indexing was last dispatched to a worker, which precedes a worker starting it. */
295295
processingQueuedAt: nullableWireDateSchema.optional(),
296+
/** When a deferred retry of a `pending` document is due. */
297+
processingDeferredUntil: nullableWireDateSchema.optional(),
296298
processingStartedAt: nullableWireDateSchema.optional(),
297299
processingCompletedAt: nullableWireDateSchema.optional(),
298300
processingError: z.string().nullable().optional(),

‎apps/sim/lib/knowledge/api/internal-route.test.ts‎

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ import {
1313
internalKnowledgeProvenanceUserId,
1414
resolveInternalKnowledgeBillingAttribution,
1515
toInternalKnowledgeConnector,
16+
toInternalKnowledgeDocument,
1617
} from '@/lib/knowledge/api/internal-route'
1718
import { resolveKnowledgeAttributedUserId } from '@/lib/knowledge/application/billing'
1819

@@ -133,6 +134,30 @@ describe('internal Knowledge execution attribution', () => {
133134
})
134135
})
135136

137+
describe('toInternalKnowledgeDocument', () => {
138+
it('presents when a deferred retry is due, so the Retry action can match the API', () => {
139+
const deferredUntil = new Date('2026-09-01T12:00:00.000Z')
140+
const presented = toInternalKnowledgeDocument({
141+
id: 'doc-1',
142+
knowledgeBaseId: 'kb-1',
143+
filename: 'a.txt',
144+
fileUrl: 'https://example.com/a.txt',
145+
fileSize: 1,
146+
mimeType: 'text/plain',
147+
chunkCount: 0,
148+
tokenCount: 0,
149+
characterCount: 0,
150+
processingStatus: 'pending',
151+
enabled: true,
152+
uploadedAt: new Date('2026-08-01T00:00:00.000Z'),
153+
processingQueuedAt: new Date('2026-08-01T00:00:00.000Z'),
154+
processingDeferredUntil: deferredUntil,
155+
})
156+
expect(presented.processingDeferredUntil).toBe(deferredUntil.toISOString())
157+
expect(presented.processingQueuedAt).toBe('2026-08-01T00:00:00.000Z')
158+
})
159+
})
160+
136161
describe('toInternalKnowledgeConnector', () => {
137162
const row = {
138163
id: 'connector-1',

‎apps/sim/lib/knowledge/api/internal-route.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,7 @@ export function toInternalKnowledgeDocument<
8989
T extends {
9090
uploadedAt: Date | string
9191
processingQueuedAt?: Date | string | null
92+
processingDeferredUntil?: Date | string | null
9293
processingStartedAt?: Date | string | null
9394
processingCompletedAt?: Date | string | null
9495
date1?: Date | string | null
@@ -99,6 +100,7 @@ export function toInternalKnowledgeDocument<
99100
...document,
100101
uploadedAt: serializeDate(document.uploadedAt),
101102
processingQueuedAt: serializeNullableDate(document.processingQueuedAt ?? null),
103+
processingDeferredUntil: serializeNullableDate(document.processingDeferredUntil ?? null),
102104
processingStartedAt: serializeNullableDate(document.processingStartedAt ?? null),
103105
processingCompletedAt: serializeNullableDate(document.processingCompletedAt ?? null),
104106
date1: serializeNullableDate(document.date1 ?? null),

‎apps/sim/lib/knowledge/connectors/member-sync-engine.test.ts‎

Lines changed: 96 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,9 @@
11
/**
22
* @vitest-environment node
33
*/
4-
import { describe, expect, it, vi } from 'vitest'
4+
import { queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing'
5+
import { DrizzleQueryError } from 'drizzle-orm/errors'
6+
import { beforeEach, describe, expect, it, vi } from 'vitest'
57

68
vi.mock('@/connectors/registry.server', () => ({ CONNECTOR_REGISTRY: {} }))
79
vi.mock('@/lib/knowledge/documents/service', () => ({
@@ -20,11 +22,13 @@ vi.mock('@/lib/billing/core/workspace-access', () => ({
2022
vi.mock('@/lib/credential-groups/availability', () => ({ isCredentialGroupsAvailable: vi.fn() }))
2123

2224
import {
25+
buildMemberSyncDatabaseRetryUpdate,
2326
buildMemberSyncFailureUpdate,
2427
deriveMemberActive,
2528
memberFailureBackoffMs,
2629
memberNextAttemptAt,
2730
nextMemberSyncTime,
31+
resolveMemberSyncFailureUpdate,
2832
shouldListFully,
2933
} from '@/lib/knowledge/connectors/member-sync-engine'
3034
import {
@@ -254,6 +258,97 @@ describe('member sync engine decisions', () => {
254258
})
255259
})
256260

261+
describe('buildMemberSyncDatabaseRetryUpdate', () => {
262+
const now = new Date('2026-09-01T12:00:00Z')
263+
const minutesAfter = (mins: number) => now.getTime() + mins * 60 * 1000
264+
265+
it('keeps the error visible without advancing the breaker', () => {
266+
const update = buildMemberSyncDatabaseRetryUpdate(
267+
now,
268+
MAX_CONSECUTIVE_FAILURES - 1,
269+
'db timeout',
270+
40
271+
)
272+
expect(update).toMatchObject({
273+
memberSyncStatus: 'error',
274+
lastMemberSyncError: 'db timeout',
275+
memberSyncConsecutiveFailures: MAX_CONSECUTIVE_FAILURES - 1,
276+
memberSyncLockToken: null,
277+
memberSyncLockLeaseAt: null,
278+
})
279+
})
280+
281+
it('climbs the failure ladder with the failed-run streak, up to its ceiling', () => {
282+
const at = (streak: number) =>
283+
buildMemberSyncDatabaseRetryUpdate(now, 0, 'db timeout', streak).nextMemberSyncAt.getTime()
284+
expect(at(1)).toBeGreaterThanOrEqual(minutesAfter(30))
285+
expect(at(1)).toBeLessThanOrEqual(minutesAfter(31))
286+
expect(at(4)).toBeGreaterThanOrEqual(minutesAfter(120))
287+
expect(at(4)).toBeLessThanOrEqual(minutesAfter(121))
288+
expect(at(500)).toBeLessThanOrEqual(minutesAfter(24 * 60 + 1))
289+
})
290+
})
291+
292+
describe('resolveMemberSyncFailureUpdate', () => {
293+
beforeEach(() => {
294+
resetDbChainMock()
295+
})
296+
297+
const failure = {
298+
connectorId: 'c-1',
299+
runId: 'run-1',
300+
previousFailures: MAX_CONSECUTIVE_FAILURES - 1,
301+
errorMessage: 'failed',
302+
}
303+
304+
it('does not disable a connector one failure from the breaker over a database timeout', async () => {
305+
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [
306+
{ status: 'failed' },
307+
{ status: 'completed' },
308+
])
309+
const timeout = new DrizzleQueryError(
310+
'select private SQL',
311+
['private'],
312+
Object.assign(new Error('canceling statement due to statement timeout'), {
313+
code: '57014',
314+
})
315+
)
316+
const update = await resolveMemberSyncFailureUpdate(timeout, failure)
317+
expect(update).toMatchObject({
318+
memberSyncStatus: 'error',
319+
memberSyncConsecutiveFailures: MAX_CONSECUTIVE_FAILURES - 1,
320+
})
321+
expect(update.nextMemberSyncAt).not.toBeNull()
322+
})
323+
324+
it('reads the members-mode run log for the streak', async () => {
325+
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [
326+
{ status: 'failed' },
327+
{ status: 'failed' },
328+
{ status: 'completed' },
329+
])
330+
const deadlock = new DrizzleQueryError(
331+
'update private SQL',
332+
['private'],
333+
Object.assign(new Error('deadlock detected'), { code: '40P01' })
334+
)
335+
const before = Date.now()
336+
const update = await resolveMemberSyncFailureUpdate(deadlock, {
337+
...failure,
338+
previousFailures: 0,
339+
})
340+
expect(update.nextMemberSyncAt!.getTime() - before).toBeGreaterThanOrEqual(90 * 60 * 1000)
341+
})
342+
343+
it('still disables at the breaker for a failure the database did not cause', async () => {
344+
const update = await resolveMemberSyncFailureUpdate(new Error('source broke'), failure)
345+
expect(update).toMatchObject({
346+
memberSyncStatus: 'disabled',
347+
memberSyncConsecutiveFailures: MAX_CONSECUTIVE_FAILURES,
348+
})
349+
})
350+
})
351+
257352
describe('nextMemberSyncTime', () => {
258353
const now = new Date('2026-09-01T12:00:00Z')
259354

‎apps/sim/lib/knowledge/connectors/member-sync-engine.ts‎

Lines changed: 79 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ import {
1111
knowledgeDocumentObservation,
1212
} from '@sim/db/schema'
1313
import { createLogger } from '@sim/logger'
14-
import { getErrorMessage, toError } from '@sim/utils/errors'
14+
import { getErrorMessage, getTransientDatabaseFailure, toError } from '@sim/utils/errors'
1515
import { generateId } from '@sim/utils/id'
1616
import { randomInt } from '@sim/utils/random'
1717
import { and, asc, eq, gt, inArray, isNull, lte, notExists, sql } from 'drizzle-orm'
@@ -65,6 +65,10 @@ import {
6565
} from '@/lib/knowledge/connectors/member-observations'
6666
import { inviteWorkspaceMembersToCredentialGroup } from '@/lib/knowledge/connectors/member-provisioning'
6767
import { runConnectorContentPass } from '@/lib/knowledge/connectors/sync-content-pass'
68+
import {
69+
countFailedRunStreak,
70+
databaseRetryDelayMs,
71+
} from '@/lib/knowledge/connectors/sync-database-retry'
6872
import {
6973
deferConnectorSync,
7074
getConnectorSyncDeferral,
@@ -327,6 +331,73 @@ export function buildMemberSyncFailureUpdate(
327331
}
328332
}
329333

334+
/**
335+
* The connector row written after the database, not the source, failed a members-mode run. The
336+
* members-mode counterpart of `buildSyncDatabaseRetryUpdate`: the breaker keeps only the source
337+
* failures already counted, while the retry climbs the failure ladder by the failed-run streak.
338+
*/
339+
export function buildMemberSyncDatabaseRetryUpdate(
340+
now: Date,
341+
previousFailures: number | null | undefined,
342+
errorMessage: string,
343+
failedRunStreak: number
344+
) {
345+
return {
346+
memberSyncStatus: 'error' as const,
347+
lastMemberSyncError: errorMessage,
348+
nextMemberSyncAt: new Date(
349+
now.getTime() + databaseRetryDelayMs(failedRunStreak, previousFailures)
350+
),
351+
memberSyncConsecutiveFailures: previousFailures ?? 0,
352+
memberSyncLockToken: null,
353+
memberSyncLockLeaseAt: null,
354+
updatedAt: now,
355+
}
356+
}
357+
358+
/**
359+
* The connector row a failed members-mode run writes. A deterministic capacity rejection waits for
360+
* an operator, a transient database failure retries without touching the breaker, and anything
361+
* else climbs the ladder toward auto-disable.
362+
*/
363+
export async function resolveMemberSyncFailureUpdate(
364+
error: unknown,
365+
failure: {
366+
connectorId: string
367+
runId: string
368+
previousFailures: number
369+
errorMessage: string
370+
retryAfterMs?: number
371+
}
372+
) {
373+
const now = new Date()
374+
if (error instanceof ConnectorSyncCapacityError) {
375+
return {
376+
memberSyncStatus: 'error' as const,
377+
lastMemberSyncError: failure.errorMessage,
378+
nextMemberSyncAt: null,
379+
memberSyncConsecutiveFailures: failure.previousFailures,
380+
memberSyncLockToken: null,
381+
memberSyncLockLeaseAt: null,
382+
updatedAt: now,
383+
}
384+
}
385+
if (getTransientDatabaseFailure(error)) {
386+
return buildMemberSyncDatabaseRetryUpdate(
387+
now,
388+
failure.previousFailures,
389+
failure.errorMessage,
390+
await countFailedRunStreak('member', failure.connectorId, failure.runId)
391+
)
392+
}
393+
return buildMemberSyncFailureUpdate(
394+
now,
395+
failure.previousFailures,
396+
failure.errorMessage,
397+
failure.retryAfterMs
398+
)
399+
}
400+
330401
/**
331402
* When a member who completed is next due: exactly one interval on, with no
332403
* jitter, so they are due whenever the connector's own (jittered) run lands.
@@ -2422,23 +2493,13 @@ export async function executeMemberSync(
24222493
logger.error('Member sync failed', { connectorId, runId, error: errorMessage, diagnostic })
24232494
try {
24242495
await failMemberSyncLog(runId, result, errorMessage)
2425-
const failureUpdate =
2426-
error instanceof ConnectorSyncCapacityError
2427-
? {
2428-
memberSyncStatus: 'error' as const,
2429-
lastMemberSyncError: errorMessage,
2430-
nextMemberSyncAt: null,
2431-
memberSyncConsecutiveFailures: connector.memberSyncConsecutiveFailures,
2432-
memberSyncLockToken: null,
2433-
memberSyncLockLeaseAt: null,
2434-
updatedAt: new Date(),
2435-
}
2436-
: buildMemberSyncFailureUpdate(
2437-
new Date(),
2438-
connector.memberSyncConsecutiveFailures,
2439-
errorMessage,
2440-
retryAfterMs
2441-
)
2496+
const failureUpdate = await resolveMemberSyncFailureUpdate(error, {
2497+
connectorId,
2498+
runId,
2499+
previousFailures: connector.memberSyncConsecutiveFailures,
2500+
errorMessage,
2501+
retryAfterMs,
2502+
})
24422503
const written = await db
24432504
.update(knowledgeConnector)
24442505
.set(failureUpdate)

0 commit comments

Comments
 (0)