Skip to content

Commit 3d6eda7

Browse files
committed
fix(knowledge): keep the deferred-retry watchdog alive through database failures
A transient database failure while checking postpones the check without spending an attempt, paced by how overdue it is and capped at the recheck interval, until a terminal bound past the recovery window; other errors, and any error past that bound, spend attempts as before.
1 parent ce14af3 commit 3d6eda7

2 files changed

Lines changed: 99 additions & 3 deletions

File tree

‎apps/sim/lib/knowledge/documents/deferred-retry-check.test.ts‎

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ vi.mock('@/lib/knowledge/documents/processing-recovery-queue', () => ({
2525

2626
import {
2727
checkDeferredDocumentRetry,
28+
DEFERRED_RETRY_CHECK_TERMINAL_MS,
2829
DEFERRED_RETRY_LOST_ERROR,
2930
DEFERRED_RETRY_RECHECK_MS,
3031
} from '@/lib/knowledge/documents/deferred-retry-check'
@@ -175,6 +176,58 @@ describe('checkDeferredDocumentRetry', () => {
175176
expect(failedWrite()).toBeUndefined()
176177
})
177178

179+
describe('database failures while checking', () => {
180+
const lockTimeout = () =>
181+
Object.assign(new Error('Failed query: private SQL'), {
182+
query: 'private SQL',
183+
cause: Object.assign(new Error('canceling statement due to lock timeout'), {
184+
code: '55P03',
185+
}),
186+
})
187+
188+
function expectDatabaseDeferral(result: unknown) {
189+
expect(result).toMatchObject({ outcome: 'deferred', consumeAttempt: false })
190+
const backoff = (result as { minimumBackoffMs: number }).minimumBackoffMs
191+
expect(backoff).toBeGreaterThan(0)
192+
expect(backoff).toBeLessThanOrEqual(DEFERRED_RETRY_RECHECK_MS)
193+
}
194+
195+
it('postpones the check without spending an attempt when the read fails', async () => {
196+
vi.setSystemTime(OVERDUE)
197+
dbChainMockFns.limit.mockRejectedValueOnce(lockTimeout())
198+
expectDatabaseDeferral(await checkDeferredDocumentRetry(PAYLOAD, context))
199+
})
200+
201+
it('postpones the check without spending an attempt when the fenced write fails', async () => {
202+
dbChainMockFns.returning.mockRejectedValueOnce(lockTimeout())
203+
expectDatabaseDeferral(await check(DEFERRED_ROW))
204+
})
205+
206+
it('backs off longer the longer the check has been overdue, up to the recheck interval', async () => {
207+
vi.setSystemTime(DEFERRED_UNTIL.getTime() + QUEUED_DISPATCH_GRACE_MS + 6 * 60 * 60 * 1000)
208+
dbChainMockFns.limit.mockRejectedValueOnce(lockTimeout())
209+
const late = (await checkDeferredDocumentRetry(PAYLOAD, context)) as {
210+
minimumBackoffMs: number
211+
}
212+
expect(late.minimumBackoffMs).toBeGreaterThanOrEqual(DEFERRED_RETRY_RECHECK_MS * 0.8)
213+
expect(late.minimumBackoffMs).toBeLessThanOrEqual(DEFERRED_RETRY_RECHECK_MS * 1.2)
214+
})
215+
216+
it('spends an attempt on an error that is not a transient database failure', async () => {
217+
vi.setSystemTime(OVERDUE)
218+
const constraint = Object.assign(new Error('duplicate key'), { code: '23505' })
219+
dbChainMockFns.limit.mockRejectedValueOnce(constraint)
220+
await expect(checkDeferredDocumentRetry(PAYLOAD, context)).rejects.toBe(constraint)
221+
})
222+
223+
it('spends attempts on a database failure once past the terminal bound, so the event ends', async () => {
224+
vi.setSystemTime(DEFERRED_UNTIL.getTime() + DEFERRED_RETRY_CHECK_TERMINAL_MS)
225+
const error = lockTimeout()
226+
dbChainMockFns.limit.mockRejectedValueOnce(error)
227+
await expect(checkDeferredDocumentRetry(PAYLOAD, context)).rejects.toBe(error)
228+
})
229+
})
230+
178231
it('rejects a payload without its document or deferral', async () => {
179232
await expect(checkDeferredDocumentRetry({ knowledgeBaseId: 'kb-1' }, context)).rejects.toThrow(
180233
'missing'

‎apps/sim/lib/knowledge/documents/deferred-retry-check.ts‎

Lines changed: 46 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,9 @@ import { db } from '@sim/db'
22
import { document } from '@sim/db/schema'
33
import { createLogger } from '@sim/logger'
44
import { toStringOrNull } from '@sim/utils/coerce'
5+
import { getTransientDatabaseFailure } from '@sim/utils/errors'
56
import { toRecord } from '@sim/utils/object'
7+
import { backoffWithJitter } from '@sim/utils/retry'
68
import { and, eq, isNotNull, isNull } from 'drizzle-orm'
79
import {
810
type DeferredOutboxHandlerResult,
@@ -23,6 +25,17 @@ const logger = createLogger('KnowledgeDeferredRetryCheck')
2325
/** How often a document whose run still looks live, or could not be looked up, is checked again. */
2426
export const DEFERRED_RETRY_RECHECK_MS = 60 * 60 * 1000
2527

28+
/**
29+
* How long after the deferral a database failure may still postpone the check without spending
30+
* an attempt. The check itself stops asking about liveness at {@link RECOVERY_WINDOW_MS}; the extra
31+
* day lets that final write outlast a slow database window too. Past it, a database failure spends
32+
* attempts like any other error, so the event always reaches completion or dead letter.
33+
*/
34+
export const DEFERRED_RETRY_CHECK_TERMINAL_MS = RECOVERY_WINDOW_MS + 24 * 60 * 60 * 1000
35+
36+
/** First delay after a database failure; later ones grow with how overdue the check is. */
37+
const DATABASE_BACKOFF_STEP_MS = 2 * 60 * 1000
38+
2639
export const DEFERRED_RETRY_LOST_ERROR =
2740
'The scheduled retry for this document did not run. Retry the document to process it again.'
2841

@@ -78,13 +91,43 @@ function isSameDeferral(
7891
* without spending an attempt. Past {@link RECOVERY_WINDOW_MS} after the deferral no retry of that
7992
* generation can still be running (a database retry is due within minutes and each run is bounded),
8093
* so the check stops asking and fails the document, which is what guarantees the event ends.
94+
*
95+
* A transient database failure while checking, most likely the same slow window that deferred the
96+
* document, postpones the check without spending an attempt, so the watchdog cannot dead-letter
97+
* during the outage it exists to outlast. That stops at {@link DEFERRED_RETRY_CHECK_TERMINAL_MS};
98+
* any other error spends an attempt as before.
8199
*/
82100
export const checkDeferredDocumentRetry: OutboxHandler<unknown> = async (
83101
rawPayload,
84102
context
85103
): Promise<DeferredOutboxHandlerResult | undefined> => {
86104
const payload = parsePayload(rawPayload)
87-
context.signal.throwIfAborted()
105+
try {
106+
return await checkDeferral(payload, context.signal)
107+
} catch (error) {
108+
context.signal.throwIfAborted()
109+
const failure = getTransientDatabaseFailure(error)
110+
const deferredUntil = Date.parse(payload.processingDeferredUntil)
111+
const now = Date.now()
112+
if (!failure || now >= deferredUntil + DEFERRED_RETRY_CHECK_TERMINAL_MS) throw error
113+
/** Paced by how long the check has been overdue, since a deferral spends no attempt to count. */
114+
const overdueMs = Math.max(0, now - deferredUntil - QUEUED_DISPATCH_GRACE_MS)
115+
return deferOutboxHandler(
116+
`Database ${failure} failure while checking the deferred retry`,
117+
backoffWithJitter(1 + Math.floor(overdueMs / DATABASE_BACKOFF_STEP_MS), null, {
118+
baseMs: DATABASE_BACKOFF_STEP_MS,
119+
maxMs: DEFERRED_RETRY_RECHECK_MS,
120+
}),
121+
false
122+
)
123+
}
124+
}
125+
126+
async function checkDeferral(
127+
payload: DeferredRetryCheckPayload,
128+
signal: AbortSignal
129+
): Promise<DeferredOutboxHandlerResult | undefined> {
130+
signal.throwIfAborted()
88131
const [row] = await db
89132
.select({
90133
...processingSnapshotColumns,
@@ -109,7 +152,7 @@ export const checkDeferredDocumentRetry: OutboxHandler<unknown> = async (
109152
return deferOutboxHandler('Deferred retry is not overdue yet', dueAt - now, false)
110153
}
111154
if (now < deferredUntil + RECOVERY_WINDOW_MS) {
112-
const { abandoned } = await inspectDocumentProcessingLiveness([row], context.signal)
155+
const { abandoned } = await inspectDocumentProcessingLiveness([row], signal)
113156
if (abandoned.length === 0) {
114157
return deferOutboxHandler(
115158
'Deferred retry may still be running',
@@ -118,7 +161,7 @@ export const checkDeferredDocumentRetry: OutboxHandler<unknown> = async (
118161
)
119162
}
120163
}
121-
context.signal.throwIfAborted()
164+
signal.throwIfAborted()
122165

123166
const failed = await db
124167
.update(document)

0 commit comments

Comments
 (0)