@@ -82,7 +82,11 @@ import { searchScopedKnowledge } from '@/lib/knowledge/application/workspace-sea
8282import { createContentSyncLease } from '@/lib/knowledge/connectors/sync-lock'
8383import { addDocument } from '@/lib/knowledge/connectors/sync-persistence'
8484import { sweepStuckDocuments } from '@/lib/knowledge/connectors/sync-primitives'
85- import { enqueueKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-outbox-event'
85+ import { DEFERRED_RETRY_LOST_ERROR } from '@/lib/knowledge/documents/deferred-retry-check'
86+ import {
87+ enqueueKnowledgeDocumentProcessing ,
88+ KNOWLEDGE_DOCUMENT_DEFERRED_RETRY_CHECK_EVENT ,
89+ } from '@/lib/knowledge/documents/processing-outbox-event'
8690import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
8791import {
8892 DOCUMENT_RECOVERY_BATCH_SIZE ,
@@ -1053,3 +1057,211 @@ describe('independent recovery of retained connector documents', () => {
10531057 }
10541058 } )
10551059} )
1060+
1061+ describe ( 'uploaded documents whose scheduled database retry is lost' , ( ) => {
1062+ const GENERATION = 'deferred-upload-generation'
1063+
1064+ async function retryChecksFor ( documentId : string ) {
1065+ return db
1066+ . select ( )
1067+ . from ( outboxEvent )
1068+ . where (
1069+ and (
1070+ eq ( outboxEvent . eventType , KNOWLEDGE_DOCUMENT_DEFERRED_RETRY_CHECK_EVENT ) ,
1071+ sql `${ outboxEvent . payload } ->>'documentId' = ${ documentId } `
1072+ )
1073+ )
1074+ }
1075+
1076+ /** Throws once, right after the claim, the way a database capacity window fails a run. */
1077+ function failNextRunAfterClaim ( error : Error ) {
1078+ return vi
1079+ . spyOn ( billingAttribution , 'assertBillingAttributionOwner' )
1080+ . mockImplementationOnce ( ( ) => {
1081+ throw error
1082+ } )
1083+ }
1084+
1085+ /** An uploaded document whose run hit a lock timeout and scheduled a retry of the same run. */
1086+ async function deferredUpload ( queuedAt : Date | null ) {
1087+ const ids = await seed ( )
1088+ const file = await failedFile ( ids )
1089+ await db
1090+ . update ( document )
1091+ . set ( {
1092+ connectorId : null ,
1093+ processingStatus : 'pending' ,
1094+ processingQueueToken : GENERATION ,
1095+ processingQueuedAt : queuedAt ,
1096+ processingCompletedAt : null ,
1097+ processingError : null ,
1098+ uploadedAt : new Date ( ) ,
1099+ } )
1100+ . where ( eq ( document . id , file . documentId ) )
1101+ const billing = await resolveSystemBillingAttribution ( ids . workspaceId )
1102+ const retryAt = new Date ( Date . now ( ) + 120_000 )
1103+ const runRetry = ( options : { onClaimed ?: ( ) => void } = { } ) =>
1104+ processDocumentAsync ( ids . knowledgeBaseId , file . documentId , file , { } , billing , GENERATION , {
1105+ processingQueueToken : GENERATION ,
1106+ chargedAtDispatch : false ,
1107+ scheduleDatabaseRetry : ( ) => retryAt ,
1108+ ...options ,
1109+ } )
1110+ const spy = failNextRunAfterClaim (
1111+ Object . assign ( new Error ( 'canceling statement due to lock timeout' ) , { code : '55P03' } )
1112+ )
1113+ try {
1114+ await expect ( runRetry ( ) ) . rejects . toMatchObject ( { code : '55P03' } )
1115+ } finally {
1116+ spy . mockRestore ( )
1117+ }
1118+ const [ check ] = await retryChecksFor ( file . documentId )
1119+ return { ids, file, billing, retryAt, check, runRetry }
1120+ }
1121+
1122+ async function runCheckAt ( eventId : string , at : number ) {
1123+ vi . useFakeTimers ( { toFake : [ 'Date' ] } )
1124+ vi . setSystemTime ( at )
1125+ try {
1126+ return await outbox . processOutboxEventById ( eventId , knowledgeDocumentProcessingOutboxHandlers )
1127+ } finally {
1128+ vi . useRealTimers ( )
1129+ }
1130+ }
1131+
1132+ const overdue = ( retryAt : Date ) => retryAt . getTime ( ) + QUEUED_DISPATCH_GRACE_MS + 60_000
1133+
1134+ it ( 'commits the deferral and its check together, due once the retry is past the grace' , async ( ) => {
1135+ const { file, retryAt, check } = await deferredUpload ( new Date ( ) )
1136+ const [ row ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1137+ expect ( row ) . toMatchObject ( { processingStatus : 'pending' , processingDeferredUntil : retryAt } )
1138+ expect ( check . availableAt ) . toEqual ( new Date ( retryAt . getTime ( ) + QUEUED_DISPATCH_GRACE_MS ) )
1139+ expect ( check . payload ) . toMatchObject ( {
1140+ documentId : file . documentId ,
1141+ processingQueueToken : GENERATION ,
1142+ processingDeferredUntil : retryAt . toISOString ( ) ,
1143+ } )
1144+ expect (
1145+ await outbox . processOutboxEventById ( check . id , knowledgeDocumentProcessingOutboxHandlers )
1146+ ) . toBe ( 'pending' )
1147+ } )
1148+
1149+ it ( 'fails the document once its scheduled retry is overdue and no run is live' , async ( ) => {
1150+ const { file, retryAt, check } = await deferredUpload ( new Date ( ) )
1151+ expect ( await runCheckAt ( check . id , overdue ( retryAt ) ) ) . toBe ( 'completed' )
1152+ const [ row ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1153+ expect ( row ) . toMatchObject ( {
1154+ processingStatus : 'failed' ,
1155+ processingError : DEFERRED_RETRY_LOST_ERROR ,
1156+ processingDeferredUntil : null ,
1157+ processingQueueToken : GENERATION ,
1158+ } )
1159+ expect ( row . processingCompletedAt ) . not . toBeNull ( )
1160+ } )
1161+
1162+ it ( 'checks again later, without spending an attempt, while the retry run is live' , async ( ) => {
1163+ const { file, retryAt, check } = await deferredUpload ( new Date ( ) )
1164+ fixture . useTrigger = true
1165+ fixture . listRuns . mockResolvedValue ( {
1166+ data : [ { id : 'run-delayed' , status : 'DELAYED' } ] ,
1167+ hasNextPage : ( ) => false ,
1168+ } )
1169+ expect ( await runCheckAt ( check . id , overdue ( retryAt ) ) ) . toBe ( 'pending' )
1170+ const [ event ] = await retryChecksFor ( file . documentId )
1171+ expect ( event . attempts ) . toBe ( 0 )
1172+ const [ row ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1173+ expect ( row ) . toMatchObject ( { processingStatus : 'pending' , processingDeferredUntil : retryAt } )
1174+ } )
1175+
1176+ it . each ( [
1177+ [ 'claims' , { processingStatus : 'processing' , processingStartedAt : new Date ( ) } ] ,
1178+ [ 'claims and re-defers' , { processingDeferredUntil : new Date ( Date . now ( ) + 600_000 ) } ] ,
1179+ ] as const ) (
1180+ 'never overwrites a retry that %s the document while the check inspects it' ,
1181+ async ( _label , change ) => {
1182+ const { file, retryAt, check } = await deferredUpload ( new Date ( ) )
1183+ fixture . useTrigger = true
1184+ fixture . listRuns . mockImplementation ( async ( ) => {
1185+ await db . update ( document ) . set ( change ) . where ( eq ( document . id , file . documentId ) )
1186+ return { data : [ ] , hasNextPage : ( ) => false }
1187+ } )
1188+ expect ( await runCheckAt ( check . id , overdue ( retryAt ) ) ) . toBe ( 'completed' )
1189+ const [ row ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1190+ expect ( row ) . toMatchObject ( change )
1191+ expect ( row . processingStatus ) . not . toBe ( 'failed' )
1192+ }
1193+ )
1194+
1195+ it ( 'schedules no check for a connector document, which the recovery sweep covers' , async ( ) => {
1196+ const ids = await seed ( )
1197+ const file = await failedFile ( ids )
1198+ await db
1199+ . update ( document )
1200+ . set ( { processingStatus : 'pending' , processingQueueToken : GENERATION } )
1201+ . where ( eq ( document . id , file . documentId ) )
1202+ const spy = failNextRunAfterClaim (
1203+ Object . assign ( new Error ( 'canceling statement due to lock timeout' ) , { code : '55P03' } )
1204+ )
1205+ try {
1206+ await expect (
1207+ processDocumentAsync (
1208+ ids . knowledgeBaseId ,
1209+ file . documentId ,
1210+ file ,
1211+ { } ,
1212+ await resolveSystemBillingAttribution ( ids . workspaceId ) ,
1213+ GENERATION ,
1214+ {
1215+ processingQueueToken : GENERATION ,
1216+ chargedAtDispatch : false ,
1217+ scheduleDatabaseRetry : ( ) => new Date ( Date . now ( ) + 120_000 ) ,
1218+ }
1219+ )
1220+ ) . rejects . toMatchObject ( { code : '55P03' } )
1221+ } finally {
1222+ spy . mockRestore ( )
1223+ }
1224+ const [ row ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1225+ expect ( row . processingStatus ) . toBe ( 'pending' )
1226+ expect ( await retryChecksFor ( file . documentId ) ) . toHaveLength ( 0 )
1227+ } )
1228+
1229+ it ( 'keeps a dispatch from claiming a deferred run that was never stamped, and the retry still claims it' , async ( ) => {
1230+ const { ids, file, billing, retryAt, runRetry } = await deferredUpload ( null )
1231+ const [ deferred ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1232+ expect ( deferred . processingQueuedAt ) . toEqual ( retryAt )
1233+
1234+ await processDocumentsWithQueue (
1235+ [
1236+ {
1237+ documentId : file . documentId ,
1238+ filename : deferred . filename ,
1239+ fileUrl : deferred . fileUrl ,
1240+ fileSize : deferred . fileSize ,
1241+ mimeType : deferred . mimeType ,
1242+ } ,
1243+ ] ,
1244+ ids . knowledgeBaseId ,
1245+ { } ,
1246+ generateId ( ) ,
1247+ billing ,
1248+ 'interactive'
1249+ )
1250+ const [ afterDispatch ] = await db . select ( ) . from ( document ) . where ( eq ( document . id , file . documentId ) )
1251+ expect ( afterDispatch ) . toMatchObject ( {
1252+ processingStatus : 'pending' ,
1253+ processingQueueToken : GENERATION ,
1254+ processingAttempts : deferred . processingAttempts ,
1255+ processingDeferredUntil : retryAt ,
1256+ } )
1257+
1258+ const onClaimed = vi . fn ( )
1259+ const spy = failNextRunAfterClaim ( new Error ( 'Synthetic failure after the retry claimed' ) )
1260+ try {
1261+ await runRetry ( { onClaimed } ) . catch ( ( ) => undefined )
1262+ } finally {
1263+ spy . mockRestore ( )
1264+ }
1265+ expect ( onClaimed ) . toHaveBeenCalledTimes ( 1 )
1266+ } )
1267+ } )
0 commit comments