11import { createLogger } from '@sim/logger'
2+ import { findCause } from '@sim/utils/errors'
23import { queue , task } from '@trigger.dev/sdk'
34import { env , envNumber } from '@/lib/core/config/env'
5+ import {
6+ type BackgroundRetryDecision ,
7+ type BackgroundRetryPolicy ,
8+ backgroundRetryAttemptCeiling ,
9+ getBackgroundRetryDecision ,
10+ getDatabaseRetryAt ,
11+ } from '@/lib/core/errors/background-retry'
412import {
513 BYOK_EMBEDDING_CREDENTIAL_REJECTION_MESSAGE ,
614 EMBEDDING_QUOTA_EXHAUSTED_MESSAGE ,
@@ -39,6 +47,46 @@ import { processDocumentAsync } from '@/lib/knowledge/documents/service'
3947const logger = createLogger ( 'TriggerKnowledgeProcessing' )
4048export { resolveQuotaContinuationDelayMs }
4149
50+ /**
51+ * Ordinary failures keep the configured short retries. A transient database failure backs off for
52+ * minutes, about an hour in total, so a slow database window does not exhaust every attempt inside
53+ * it and leave an uploaded document failed for good.
54+ */
55+ export const DOCUMENT_PROCESSING_RETRY_POLICY : BackgroundRetryPolicy = {
56+ maxAttempts : envNumber ( env . KB_CONFIG_MAX_ATTEMPTS , 3 ) ,
57+ database : { maxAttempts : 6 , baseDelayMs : 2 * 60 * 1000 , maxDelayMs : 30 * 60 * 1000 } ,
58+ }
59+
60+ /**
61+ * A database failure whose next attempt is already scheduled, and recorded on the document as
62+ * `pending` until {@link retryAt}. The message names only the database code; the driver error
63+ * stays in `cause`, which the task runner does not record.
64+ */
65+ export class DocumentProcessingDatabaseRetryError extends Error {
66+ constructor (
67+ message : string ,
68+ readonly retryAt : Date ,
69+ options : { cause : unknown }
70+ ) {
71+ super ( message , options )
72+ this . name = 'DocumentProcessingDatabaseRetryError'
73+ }
74+ }
75+
76+ /** The `catchError` decision for `knowledge-process-document` after `attempt` (1-based) failed. */
77+ export function getDocumentProcessingRetry (
78+ error : unknown ,
79+ attempt : number
80+ ) : BackgroundRetryDecision {
81+ const scheduled = findCause (
82+ error ,
83+ ( value ) : value is DocumentProcessingDatabaseRetryError =>
84+ value instanceof DocumentProcessingDatabaseRetryError
85+ )
86+ if ( scheduled ) return { retryAt : scheduled . retryAt }
87+ return getBackgroundRetryDecision ( error , attempt , DOCUMENT_PROCESSING_RETRY_POLICY )
88+ }
89+
4290export async function runDocumentProcessing (
4391 rawPayload : DocumentProcessingPayload ,
4492 attemptNumber = 1
@@ -56,6 +104,8 @@ export async function runDocumentProcessing(
56104 payload . processingSliceCount === undefined
57105
58106 logger . info ( `[${ requestId } ] Starting Trigger.dev processing for document: ${ docData . filename } ` )
107+ /** Set from the service's callback, so control-flow narrowing cannot see it change. */
108+ let databaseRetryAt = null as Date | null
59109
60110 try {
61111 const result = await processDocumentAsync (
@@ -87,6 +137,14 @@ export async function runDocumentProcessing(
87137 : { quotaContinuationExhausted : true } ) ,
88138 scheduleProviderContinuation : ( error ) =>
89139 scheduleDocumentProcessingProviderContinuation ( payload , error , true , chargedAtDispatch ) ,
140+ scheduleDatabaseRetry : ( error ) => {
141+ databaseRetryAt = getDatabaseRetryAt (
142+ error ,
143+ attemptNumber ,
144+ DOCUMENT_PROCESSING_RETRY_POLICY
145+ )
146+ return databaseRetryAt
147+ } ,
90148 }
91149 )
92150
@@ -100,6 +158,20 @@ export async function runDocumentProcessing(
100158 processingTime : Date . now ( ) - startedAt ,
101159 }
102160 } catch ( error ) {
161+ if ( databaseRetryAt ) {
162+ const diagnostic = getConnectorFailureDiagnostic ( error )
163+ logger . warn ( `[${ requestId } ] Document processing will retry after a database failure` , {
164+ documentId,
165+ diagnostic,
166+ attempt : attemptNumber ,
167+ retryAt : databaseRetryAt . toISOString ( ) ,
168+ } )
169+ throw new DocumentProcessingDatabaseRetryError (
170+ diagnostic ?. message ?? 'Database request failed.' ,
171+ databaseRetryAt ,
172+ { cause : error }
173+ )
174+ }
103175 const providerDeferral = getProviderCapacityDeferral ( error )
104176 if ( providerDeferral || error instanceof ProviderCapacityContinuationExhaustedError ) {
105177 const outcome =
@@ -205,7 +277,8 @@ export async function runDocumentProcessing(
205277 `[${ requestId } ] Failed to process document: ${ docData . filename } ` ,
206278 diagnostic ?? error
207279 )
208- if ( diagnostic ?. category === 'database' ) throw new Error ( diagnostic . message )
280+ /** Trigger records the thrown message and stack, never `cause`; Drizzle's message carries SQL. */
281+ if ( diagnostic ?. category === 'database' ) throw new Error ( diagnostic . message , { cause : error } )
209282 throw error
210283 }
211284}
@@ -253,7 +326,13 @@ export const processDocument = task({
253326 */
254327 machine : 'medium-2x' ,
255328 retry : {
256- maxAttempts : envNumber ( env . KB_CONFIG_MAX_ATTEMPTS , 3 ) ,
329+ /**
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.
334+ */
335+ maxAttempts : backgroundRetryAttemptCeiling ( DOCUMENT_PROCESSING_RETRY_POLICY ) ,
257336 factor : envNumber ( env . KB_CONFIG_RETRY_FACTOR , 2 ) ,
258337 minTimeoutInMs : envNumber ( env . KB_CONFIG_MIN_TIMEOUT , 1000 ) ,
259338 maxTimeoutInMs : envNumber ( env . KB_CONFIG_MAX_TIMEOUT , 10000 ) ,
@@ -272,4 +351,5 @@ export const processDocument = task({
272351 queue : interactiveProcessingQueue ,
273352 run : ( payload : DocumentProcessingPayload , { ctx } ) =>
274353 runDocumentProcessing ( payload , ctx . attempt . number ) ,
354+ catchError : async ( { error, ctx } ) => getDocumentProcessingRetry ( error , ctx . attempt . number ) ,
275355} )
0 commit comments