Skip to content

Commit 778d445

Browse files
fix(knowledge): make indexing and connector recovery durable (#7670)
* fix(knowledge): make indexing and connector recovery durable * fix(knowledge): expose safe recovery diagnostics
1 parent a77b1a8 commit 778d445

58 files changed

Lines changed: 29238 additions & 535 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: 30 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,11 +10,14 @@ import { enterpriseIssuanceOutboxHandlers } from '@/lib/billing/enterprise-provi
1010
import { membershipBillingOutboxHandlers } from '@/lib/billing/organizations/membership-reconciliation'
1111
import { billingOutboxHandlers } from '@/lib/billing/webhooks/outbox-handlers'
1212
import { processOutboxEvents } from '@/lib/core/outbox/service'
13+
import { DeadlineExceededError } from '@/lib/core/utils/deadline'
1314
import { generateRequestId } from '@/lib/core/utils/request'
1415
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
1516
import { directGrantOutboxHandlers } from '@/lib/invitations/direct-grant'
1617
import { slackSearchOutboxHandlers } from '@/lib/knowledge/application/slack-search/outbox'
18+
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
1719
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
20+
import { recoverKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-recovery'
1821
import { organizationResourceCleanupOutboxHandlers } from '@/lib/organizations/resource-cleanup'
1922
import { workspaceFileLiveDocOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-live-doc-outbox'
2023
import { workspaceFileStorageCleanupOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox'
@@ -53,12 +56,31 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
5356
return authError
5457
}
5558

59+
const startedAt = Date.now()
5660
const result = await processOutboxEvents(handlers, {
5761
batchSize: 500,
58-
maxRuntimeMs: 790_000,
62+
maxRuntimeMs: 760_000,
5963
minRemainingMs: 95_000,
6064
})
6165

66+
let recoveredDocuments = 0
67+
try {
68+
if (Date.now() - startedAt < 770_000) {
69+
recoveredDocuments = await recoverKnowledgeDocumentProcessing()
70+
}
71+
} catch (error) {
72+
logger.error('Stored document recovery failed', {
73+
requestId,
74+
error: getConnectorFailureDiagnostic(error) ?? {
75+
category: error instanceof DeadlineExceededError ? 'deadline' : 'internal',
76+
message:
77+
error instanceof DeadlineExceededError
78+
? error.message
79+
: 'Unexpected stored-document recovery failure',
80+
},
81+
})
82+
}
83+
6284
// Reap fork background-work rows stuck `processing` past their TTL (worker crash /
6385
// restart has no in-task hook). Independent of the outbox; a failure here must not
6486
// fail the outbox run, so it's guarded separately.
@@ -69,13 +91,19 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
6991
logger.error('Background-work reap failed', { requestId, error: toError(error).message })
7092
}
7193

72-
logger.info('Outbox processing completed', { requestId, ...result, reapedBackgroundWork })
94+
logger.info('Outbox processing completed', {
95+
requestId,
96+
...result,
97+
reapedBackgroundWork,
98+
recoveredDocuments,
99+
})
73100

74101
return NextResponse.json({
75102
success: true,
76103
requestId,
77104
result,
78105
reapedBackgroundWork,
106+
recoveredDocuments,
79107
})
80108
} catch (error) {
81109
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)