Skip to content

Commit 74eec09

Browse files
committed
fix(knowledge): make indexing and connector recovery durable
1 parent e508098 commit 74eec09

58 files changed

Lines changed: 29206 additions & 532 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

apps/sim/app/api/webhooks/outbox/process/route.ts

Lines changed: 19 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
1515
import { directGrantOutboxHandlers } from '@/lib/invitations/direct-grant'
1616
import { slackSearchOutboxHandlers } from '@/lib/knowledge/application/slack-search/outbox'
1717
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
18+
import { recoverKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-recovery'
1819
import { organizationResourceCleanupOutboxHandlers } from '@/lib/organizations/resource-cleanup'
1920
import { workspaceFileLiveDocOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-live-doc-outbox'
2021
import { workspaceFileStorageCleanupOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox'
@@ -53,12 +54,22 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
5354
return authError
5455
}
5556

57+
const startedAt = Date.now()
5658
const result = await processOutboxEvents(handlers, {
5759
batchSize: 500,
58-
maxRuntimeMs: 790_000,
60+
maxRuntimeMs: 760_000,
5961
minRemainingMs: 95_000,
6062
})
6163

64+
let recoveredDocuments = 0
65+
try {
66+
if (Date.now() - startedAt < 770_000) {
67+
recoveredDocuments = await recoverKnowledgeDocumentProcessing()
68+
}
69+
} catch {
70+
logger.error('Stored document recovery failed', { requestId })
71+
}
72+
6273
// Reap fork background-work rows stuck `processing` past their TTL (worker crash /
6374
// restart has no in-task hook). Independent of the outbox; a failure here must not
6475
// fail the outbox run, so it's guarded separately.
@@ -69,13 +80,19 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
6980
logger.error('Background-work reap failed', { requestId, error: toError(error).message })
7081
}
7182

72-
logger.info('Outbox processing completed', { requestId, ...result, reapedBackgroundWork })
83+
logger.info('Outbox processing completed', {
84+
requestId,
85+
...result,
86+
reapedBackgroundWork,
87+
recoveredDocuments,
88+
})
7389

7490
return NextResponse.json({
7591
success: true,
7692
requestId,
7793
result,
7894
reapedBackgroundWork,
95+
recoveredDocuments,
7996
})
8097
} catch (error) {
8198
logger.error('Outbox processing failed', { requestId, error: toError(error).message })

apps/sim/background/knowledge-connector-member-sync.ts

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,12 +13,19 @@ import { MEMBER_SYNC_MAX_DURATION_SECONDS } from '@/lib/knowledge/connectors/syn
1313

1414
const logger = createLogger('TriggerKnowledgeConnectorMemberSync')
1515

16-
export type MemberSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed'
16+
export type MemberSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed' | 'deferred'
1717

1818
/** A run is partial when any member or document failed; skipped and failed mirror the content task. */
1919
export function classifyMemberSyncResult(result: MemberSyncResult): MemberSyncTaskOutcome {
2020
if (result.skipReason) return 'skipped'
2121
if (result.error) return 'failed'
22+
if (
23+
result.deferred &&
24+
result.docsFailed === 0 &&
25+
result.processingDispatch.failed === 0 &&
26+
result.membersFailed === 0
27+
)
28+
return 'deferred'
2229
if (
2330
result.listingIncomplete ||
2431
result.membersIncomplete > 0 ||
@@ -49,6 +56,7 @@ export async function executeMemberSyncJob(payload: unknown) {
4956
logger.info(`[${requestId}] Member sync completed`, {
5057
connectorId,
5158
outcome,
59+
deferred: result.deferred,
5260
membersClaimed: result.membersClaimed,
5361
membersCompleted: result.membersCompleted,
5462
membersIncomplete: result.membersIncomplete,

apps/sim/background/knowledge-connector-sync.test.ts

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,34 @@ describe('knowledge connector sync worker', () => {
140140
await expect(run).rejects.toThrow('Connector sync partially failed')
141141
})
142142

143+
it('completes a durably scheduled capacity wait while preserving existing source failures', async () => {
144+
mockAssertConnectorSyncPayload.mockReturnValue({
145+
connectorId: 'connector-1',
146+
requestId: 'request-1',
147+
billingAttribution: BILLING_ATTRIBUTION,
148+
})
149+
const base = await mockExecuteSync()
150+
const waiting = {
151+
...base,
152+
listingIncomplete: true,
153+
deferred: {
154+
reason: 'admission_timeout',
155+
providerId: 'github-rest',
156+
nextSyncAt: '2026-09-01T00:00:00Z',
157+
},
158+
}
159+
mockExecuteSync.mockResolvedValue(waiting)
160+
expect(await executeConnectorSyncJob({})).toMatchObject({
161+
outcome: 'deferred',
162+
success: false,
163+
deferred: waiting.deferred,
164+
})
165+
mockExecuteSync.mockResolvedValue({ ...waiting, docsFailed: 1 })
166+
await expect(executeConnectorSyncJob({})).rejects.toThrow('partially failed')
167+
mockExecuteSync.mockResolvedValue({ ...waiting, error: 'Retry persistence failed' })
168+
await expect(executeConnectorSyncJob({})).rejects.toThrow('Retry persistence failed')
169+
})
170+
143171
it('does not turn intentionally skipped source files into a task failure', () => {
144172
expect(
145173
classifyConnectorSyncResult({

apps/sim/background/knowledge-connector-sync.ts

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ import type { SyncResult } from '@/connectors/types'
1010

1111
const logger = createLogger('TriggerKnowledgeConnectorSync')
1212

13-
export type ConnectorSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed'
13+
export type ConnectorSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed' | 'deferred'
1414

1515
/**
1616
* Separates source-sync failures from expected queue/lock no-ops. Intentional
@@ -20,6 +20,8 @@ export type ConnectorSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'fa
2020
export function classifyConnectorSyncResult(result: SyncResult): ConnectorSyncTaskOutcome {
2121
if (result.skipReason) return 'skipped'
2222
if (result.error) return 'failed'
23+
if (result.deferred && result.docsFailed === 0 && result.processingDispatch.failed === 0)
24+
return 'deferred'
2325
if (result.listingIncomplete || result.docsFailed > 0 || result.processingDispatch.failed > 0)
2426
return 'partial'
2527
return 'completed'
@@ -61,6 +63,7 @@ export async function executeConnectorSyncJob(payload: unknown) {
6163
logger.info(`[${requestId}] Connector sync completed`, {
6264
connectorId,
6365
outcome: classifyConnectorSyncResult(result),
66+
deferred: result.deferred,
6467
added: result.docsAdded,
6568
updated: result.docsUpdated,
6669
deleted: result.docsDeleted,

0 commit comments

Comments
 (0)