Skip to content

Commit 91c0228

Browse files
committed
fix(jobs): progress-aware connector database retries, stuck-sync signal, strict query shapes
- the database-failure streak counts only failed runs that made no progress; a run that wrote documents (or completed a member) retries in minutes - ten zero-progress database failures in a row log an alertable error; the connector is never disabled for them - the classifier's query check matches only Drizzle's query error and the postgres.js query error shapes, so a client error carrying its own query stays a source failure
1 parent 3d6eda7 commit 91c0228

8 files changed

Lines changed: 393 additions & 100 deletions

File tree

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

Lines changed: 33 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -278,14 +278,10 @@ describe('member sync engine decisions', () => {
278278
})
279279
})
280280

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))
281+
it('schedules the next run after the resolved retry delay', () => {
282+
expect(
283+
buildMemberSyncDatabaseRetryUpdate(now, 0, 'db timeout', 120 * 60 * 1000).nextMemberSyncAt
284+
).toEqual(new Date(minutesAfter(120)))
289285
})
290286
})
291287

@@ -299,13 +295,23 @@ describe('member sync engine decisions', () => {
299295
runId: 'run-1',
300296
previousFailures: MAX_CONSECUTIVE_FAILURES - 1,
301297
errorMessage: 'failed',
298+
madeProgress: false,
302299
}
300+
const run = (status: string, membersCompleted = 0) => ({
301+
status,
302+
membersCompleted,
303+
docsAdded: 0,
304+
docsUpdated: 0,
305+
})
306+
const deadlock = () =>
307+
new DrizzleQueryError(
308+
'update private SQL',
309+
['private'],
310+
Object.assign(new Error('deadlock detected'), { code: '40P01' })
311+
)
303312

304313
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-
])
314+
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [run('failed'), run('completed')])
309315
const timeout = new DrizzleQueryError(
310316
'select private SQL',
311317
['private'],
@@ -323,23 +329,29 @@ describe('member sync engine decisions', () => {
323329

324330
it('reads the members-mode run log for the streak', async () => {
325331
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [
326-
{ status: 'failed' },
327-
{ status: 'failed' },
328-
{ status: 'completed' },
332+
run('failed'),
333+
run('failed'),
334+
run('completed'),
329335
])
330-
const deadlock = new DrizzleQueryError(
331-
'update private SQL',
332-
['private'],
333-
Object.assign(new Error('deadlock detected'), { code: '40P01' })
334-
)
335336
const before = Date.now()
336-
const update = await resolveMemberSyncFailureUpdate(deadlock, {
337+
const update = await resolveMemberSyncFailureUpdate(deadlock(), {
337338
...failure,
338339
previousFailures: 0,
339340
})
340341
expect(update.nextMemberSyncAt!.getTime() - before).toBeGreaterThanOrEqual(90 * 60 * 1000)
341342
})
342343

344+
it('retries within minutes after a run that completed members before the database failed', async () => {
345+
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [run('failed'), run('failed')])
346+
const before = Date.now()
347+
const update = await resolveMemberSyncFailureUpdate(deadlock(), {
348+
...failure,
349+
madeProgress: true,
350+
})
351+
expect(update.nextMemberSyncAt!.getTime() - before).toBeLessThanOrEqual(5 * 60 * 1000)
352+
expect(update.memberSyncConsecutiveFailures).toBe(MAX_CONSECUTIVE_FAILURES - 1)
353+
})
354+
343355
it('still disables at the breaker for a failure the database did not cause', async () => {
344356
const update = await resolveMemberSyncFailureUpdate(new Error('source broke'), failure)
345357
expect(update).toMatchObject({

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

Lines changed: 14 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -65,10 +65,7 @@ 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'
68+
import { resolveDatabaseRetryDelayMs } from '@/lib/knowledge/connectors/sync-database-retry'
7269
import {
7370
deferConnectorSync,
7471
getConnectorSyncDeferral,
@@ -334,20 +331,18 @@ export function buildMemberSyncFailureUpdate(
334331
/**
335332
* The connector row written after the database, not the source, failed a members-mode run. The
336333
* 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.
334+
* failures already counted, and the retry waits the delay `resolveDatabaseRetryDelayMs` chose.
338335
*/
339336
export function buildMemberSyncDatabaseRetryUpdate(
340337
now: Date,
341338
previousFailures: number | null | undefined,
342339
errorMessage: string,
343-
failedRunStreak: number
340+
retryDelayMs: number
344341
) {
345342
return {
346343
memberSyncStatus: 'error' as const,
347344
lastMemberSyncError: errorMessage,
348-
nextMemberSyncAt: new Date(
349-
now.getTime() + databaseRetryDelayMs(failedRunStreak, previousFailures)
350-
),
345+
nextMemberSyncAt: new Date(now.getTime() + retryDelayMs),
351346
memberSyncConsecutiveFailures: previousFailures ?? 0,
352347
memberSyncLockToken: null,
353348
memberSyncLockLeaseAt: null,
@@ -368,6 +363,8 @@ export async function resolveMemberSyncFailureUpdate(
368363
previousFailures: number
369364
errorMessage: string
370365
retryAfterMs?: number
366+
/** Whether the run completed a member or wrote documents before it failed. */
367+
madeProgress: boolean
371368
}
372369
) {
373370
const now = new Date()
@@ -387,7 +384,13 @@ export async function resolveMemberSyncFailureUpdate(
387384
now,
388385
failure.previousFailures,
389386
failure.errorMessage,
390-
await countFailedRunStreak('member', failure.connectorId, failure.runId)
387+
await resolveDatabaseRetryDelayMs({
388+
kind: 'member',
389+
connectorId: failure.connectorId,
390+
runId: failure.runId,
391+
previousFailures: failure.previousFailures,
392+
madeProgress: failure.madeProgress,
393+
})
391394
)
392395
}
393396
return buildMemberSyncFailureUpdate(
@@ -2499,6 +2502,7 @@ export async function executeMemberSync(
24992502
previousFailures: connector.memberSyncConsecutiveFailures,
25002503
errorMessage,
25012504
retryAfterMs,
2505+
madeProgress: result.membersCompleted + result.docsAdded + result.docsUpdated > 0,
25022506
})
25032507
const written = await db
25042508
.update(knowledgeConnector)

‎apps/sim/lib/knowledge/connectors/sync-database-retry.test.ts‎

Lines changed: 133 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -8,50 +8,109 @@ import {
88
resetDbChainMock,
99
schemaMock,
1010
} from '@sim/testing'
11-
import { beforeEach, describe, expect, it } from 'vitest'
11+
import { beforeEach, describe, expect, it, vi } from 'vitest'
12+
13+
const { mockLogError } = vi.hoisted(() => ({ mockLogError: vi.fn() }))
14+
vi.mock('@sim/logger', async () => {
15+
const { createMockLogger } = await import('@sim/testing/mocks/logger.mock')
16+
return { createLogger: () => ({ ...createMockLogger(), error: mockLogError }) }
17+
})
18+
1219
import {
13-
countFailedRunStreak,
20+
countZeroProgressFailedRuns,
21+
DATABASE_FAILURE_ALERT_STREAK,
22+
DATABASE_RETRY_AFTER_PROGRESS_MS,
1423
databaseRetryDelayMs,
24+
resolveDatabaseRetryDelayMs,
1525
} from '@/lib/knowledge/connectors/sync-database-retry'
1626
import {
1727
CONNECTOR_FAILURE_BACKOFF_CAP_MINUTES,
1828
MAX_CONSECUTIVE_FAILURES,
1929
} from '@/lib/knowledge/connectors/sync-limits'
2030

2131
const MINUTE = 60 * 1000
22-
const runs = (...statuses: string[]) => statuses.map((status) => ({ status }))
32+
const NO_WRITES = { docsAdded: 0, docsUpdated: 0, docsDeleted: 0 }
33+
const contentRun = (status: string, writes: Partial<typeof NO_WRITES> = {}) => ({
34+
status,
35+
...NO_WRITES,
36+
...writes,
37+
})
38+
const memberRun = (status: string, writes: Record<string, number> = {}) => ({
39+
status,
40+
membersCompleted: 0,
41+
docsAdded: 0,
42+
docsUpdated: 0,
43+
...writes,
44+
})
2345

24-
describe('countFailedRunStreak', () => {
46+
describe('countZeroProgressFailedRuns', () => {
2547
beforeEach(() => {
2648
resetDbChainMock()
2749
})
2850

2951
it('counts this run plus the failed runs before it, up to the last one that did not fail', async () => {
30-
queueTableRows(
31-
schemaMock.knowledgeConnectorSyncLog,
32-
runs('failed', 'failed', 'completed', 'failed')
33-
)
34-
expect(await countFailedRunStreak('content', 'c-1', 'run-1')).toBe(3)
52+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [
53+
contentRun('failed'),
54+
contentRun('failed'),
55+
contentRun('completed'),
56+
contentRun('failed'),
57+
])
58+
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(3)
59+
})
60+
61+
it.each([
62+
['added', { docsAdded: 5 }],
63+
['updated', { docsUpdated: 1 }],
64+
['deleted', { docsDeleted: 2 }],
65+
])('ends the streak at a failed run that %s documents', async (_label, writes) => {
66+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [
67+
contentRun('failed'),
68+
contentRun('failed', writes),
69+
contentRun('failed'),
70+
])
71+
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(2)
3572
})
3673

3774
it('counts only this run after a success', async () => {
38-
queueTableRows(schemaMock.knowledgeConnectorSyncLog, runs('completed', 'failed'))
39-
expect(await countFailedRunStreak('content', 'c-1', 'run-1')).toBe(1)
75+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [
76+
contentRun('completed'),
77+
contentRun('failed'),
78+
])
79+
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(1)
4080
})
4181

4282
it('counts an unbroken history in full', async () => {
43-
queueTableRows(schemaMock.knowledgeConnectorSyncLog, runs('failed', 'failed'))
44-
expect(await countFailedRunStreak('content', 'c-1', 'run-1')).toBe(3)
83+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [
84+
contentRun('failed'),
85+
contentRun('failed'),
86+
])
87+
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(3)
4588
})
4689

4790
it('reads the members-mode run log for a members-mode run', async () => {
48-
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, runs('failed', 'started'))
49-
expect(await countFailedRunStreak('member', 'c-1', 'run-1')).toBe(2)
91+
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [
92+
memberRun('failed'),
93+
memberRun('started'),
94+
])
95+
expect(await countZeroProgressFailedRuns('member', 'c-1', 'run-1')).toBe(2)
96+
})
97+
98+
it.each([
99+
['completed a member', { membersCompleted: 1 }],
100+
['added documents', { docsAdded: 3 }],
101+
['updated documents', { docsUpdated: 1 }],
102+
])('ends a members-mode streak at a failed run that %s', async (_label, writes) => {
103+
queueTableRows(schemaMock.knowledgeConnectorMemberSyncLog, [
104+
memberRun('failed'),
105+
memberRun('failed', writes),
106+
memberRun('failed'),
107+
])
108+
expect(await countZeroProgressFailedRuns('member', 'c-1', 'run-1')).toBe(2)
50109
})
51110

52111
it('excludes the current run and reads only as far back as the ladder climbs', async () => {
53112
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [])
54-
await countFailedRunStreak('content', 'c-1', 'run-1')
113+
await countZeroProgressFailedRuns('content', 'c-1', 'run-1')
55114
const where = dbChainMockFns.where.mock.calls.at(-1)?.[0]
56115
expect(
57116
hasMockCondition(
@@ -69,7 +128,7 @@ describe('countFailedRunStreak', () => {
69128

70129
it('falls back to this run alone when the history cannot be read', async () => {
71130
dbChainMockFns.limit.mockRejectedValueOnce(new Error('canceling statement'))
72-
expect(await countFailedRunStreak('content', 'c-1', 'run-1')).toBe(1)
131+
expect(await countZeroProgressFailedRuns('content', 'c-1', 'run-1')).toBe(1)
73132
})
74133
})
75134

@@ -96,3 +155,60 @@ describe('databaseRetryDelayMs', () => {
96155
expect(delay).toBeGreaterThanOrEqual(MAX_CONSECUTIVE_FAILURES * 30 * MINUTE)
97156
})
98157
})
158+
159+
describe('resolveDatabaseRetryDelayMs', () => {
160+
beforeEach(() => {
161+
vi.clearAllMocks()
162+
resetDbChainMock()
163+
})
164+
165+
const retry = {
166+
kind: 'content' as const,
167+
connectorId: 'c-1',
168+
runId: 'run-1',
169+
previousFailures: 0,
170+
}
171+
172+
it('retries shortly after a run that made progress, without reading the streak', async () => {
173+
queueTableRows(
174+
schemaMock.knowledgeConnectorSyncLog,
175+
Array.from({ length: 20 }, () => contentRun('failed'))
176+
)
177+
const delay = await resolveDatabaseRetryDelayMs({ ...retry, madeProgress: true })
178+
expect(delay).toBeGreaterThanOrEqual(DATABASE_RETRY_AFTER_PROGRESS_MS)
179+
expect(delay).toBeLessThanOrEqual(DATABASE_RETRY_AFTER_PROGRESS_MS + MINUTE)
180+
expect(dbChainMockFns.select).not.toHaveBeenCalled()
181+
})
182+
183+
it('climbs the ladder by the zero-progress streak for a run that made none', async () => {
184+
queueTableRows(schemaMock.knowledgeConnectorSyncLog, [
185+
contentRun('failed'),
186+
contentRun('failed'),
187+
contentRun('completed'),
188+
])
189+
const delay = await resolveDatabaseRetryDelayMs({ ...retry, madeProgress: false })
190+
expect(delay).toBeGreaterThanOrEqual(90 * MINUTE)
191+
expect(delay).toBeLessThanOrEqual(91 * MINUTE)
192+
})
193+
194+
it('reports a streak that reaches the alert threshold at error level, without disabling', async () => {
195+
queueTableRows(
196+
schemaMock.knowledgeConnectorSyncLog,
197+
Array.from({ length: DATABASE_FAILURE_ALERT_STREAK - 1 }, () => contentRun('failed'))
198+
)
199+
await resolveDatabaseRetryDelayMs({ ...retry, madeProgress: false })
200+
expect(mockLogError).toHaveBeenCalledWith(
201+
'Connector sync keeps failing on the database without progress',
202+
{ connectorId: 'c-1', kind: 'content', zeroProgressFailedRuns: DATABASE_FAILURE_ALERT_STREAK }
203+
)
204+
})
205+
206+
it('stays quiet below the alert threshold', async () => {
207+
queueTableRows(
208+
schemaMock.knowledgeConnectorSyncLog,
209+
Array.from({ length: DATABASE_FAILURE_ALERT_STREAK - 2 }, () => contentRun('failed'))
210+
)
211+
await resolveDatabaseRetryDelayMs({ ...retry, madeProgress: false })
212+
expect(mockLogError).not.toHaveBeenCalled()
213+
})
214+
})

0 commit comments

Comments
 (0)