Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 30 additions & 2 deletions apps/sim/app/api/webhooks/outbox/process/route.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,14 @@ import { enterpriseIssuanceOutboxHandlers } from '@/lib/billing/enterprise-provi
import { membershipBillingOutboxHandlers } from '@/lib/billing/organizations/membership-reconciliation'
import { billingOutboxHandlers } from '@/lib/billing/webhooks/outbox-handlers'
import { processOutboxEvents } from '@/lib/core/outbox/service'
import { DeadlineExceededError } from '@/lib/core/utils/deadline'
import { generateRequestId } from '@/lib/core/utils/request'
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
import { directGrantOutboxHandlers } from '@/lib/invitations/direct-grant'
import { slackSearchOutboxHandlers } from '@/lib/knowledge/application/slack-search/outbox'
import { getConnectorFailureDiagnostic } from '@/lib/knowledge/connectors/connector-error'
import { knowledgeDocumentProcessingOutboxHandlers } from '@/lib/knowledge/documents/processing-outbox-handler'
import { recoverKnowledgeDocumentProcessing } from '@/lib/knowledge/documents/processing-recovery'
import { organizationResourceCleanupOutboxHandlers } from '@/lib/organizations/resource-cleanup'
import { workspaceFileLiveDocOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-live-doc-outbox'
import { workspaceFileStorageCleanupOutboxHandlers } from '@/lib/uploads/contexts/workspace/workspace-file-storage-cleanup-outbox'
Expand Down Expand Up @@ -53,12 +56,31 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
return authError
}

const startedAt = Date.now()
const result = await processOutboxEvents(handlers, {
batchSize: 500,
maxRuntimeMs: 790_000,
maxRuntimeMs: 760_000,
minRemainingMs: 95_000,
})

let recoveredDocuments = 0
try {
if (Date.now() - startedAt < 770_000) {
recoveredDocuments = await recoverKnowledgeDocumentProcessing()
}
} catch (error) {
logger.error('Stored document recovery failed', {
requestId,
error: getConnectorFailureDiagnostic(error) ?? {
category: error instanceof DeadlineExceededError ? 'deadline' : 'internal',
message:
error instanceof DeadlineExceededError
? error.message
: 'Unexpected stored-document recovery failure',
},
})
}

// Reap fork background-work rows stuck `processing` past their TTL (worker crash /
// restart has no in-task hook). Independent of the outbox; a failure here must not
// fail the outbox run, so it's guarded separately.
Expand All @@ -69,13 +91,19 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
logger.error('Background-work reap failed', { requestId, error: toError(error).message })
}

logger.info('Outbox processing completed', { requestId, ...result, reapedBackgroundWork })
logger.info('Outbox processing completed', {
requestId,
...result,
reapedBackgroundWork,
recoveredDocuments,
})

return NextResponse.json({
success: true,
requestId,
result,
reapedBackgroundWork,
recoveredDocuments,
})
} catch (error) {
logger.error('Outbox processing failed', { requestId, error: toError(error).message })
Expand Down
10 changes: 9 additions & 1 deletion apps/sim/background/knowledge-connector-member-sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,12 +13,19 @@ import { MEMBER_SYNC_MAX_DURATION_SECONDS } from '@/lib/knowledge/connectors/syn

const logger = createLogger('TriggerKnowledgeConnectorMemberSync')

export type MemberSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed'
export type MemberSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed' | 'deferred'

/** A run is partial when any member or document failed; skipped and failed mirror the content task. */
export function classifyMemberSyncResult(result: MemberSyncResult): MemberSyncTaskOutcome {
if (result.skipReason) return 'skipped'
if (result.error) return 'failed'
if (
result.deferred &&
result.docsFailed === 0 &&
result.processingDispatch.failed === 0 &&
result.membersFailed === 0
)
return 'deferred'
if (
result.listingIncomplete ||
result.membersIncomplete > 0 ||
Expand Down Expand Up @@ -49,6 +56,7 @@ export async function executeMemberSyncJob(payload: unknown) {
logger.info(`[${requestId}] Member sync completed`, {
connectorId,
outcome,
deferred: result.deferred,
membersClaimed: result.membersClaimed,
membersCompleted: result.membersCompleted,
membersIncomplete: result.membersIncomplete,
Expand Down
28 changes: 28 additions & 0 deletions apps/sim/background/knowledge-connector-sync.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,34 @@ describe('knowledge connector sync worker', () => {
await expect(run).rejects.toThrow('Connector sync partially failed')
})

it('completes a durably scheduled capacity wait while preserving existing source failures', async () => {
mockAssertConnectorSyncPayload.mockReturnValue({
connectorId: 'connector-1',
requestId: 'request-1',
billingAttribution: BILLING_ATTRIBUTION,
})
const base = await mockExecuteSync()
const waiting = {
...base,
listingIncomplete: true,
deferred: {
reason: 'admission_timeout',
providerId: 'github-rest',
nextSyncAt: '2026-09-01T00:00:00Z',
},
}
mockExecuteSync.mockResolvedValue(waiting)
expect(await executeConnectorSyncJob({})).toMatchObject({
outcome: 'deferred',
success: false,
deferred: waiting.deferred,
})
mockExecuteSync.mockResolvedValue({ ...waiting, docsFailed: 1 })
await expect(executeConnectorSyncJob({})).rejects.toThrow('partially failed')
mockExecuteSync.mockResolvedValue({ ...waiting, error: 'Retry persistence failed' })
await expect(executeConnectorSyncJob({})).rejects.toThrow('Retry persistence failed')
})

it('does not turn intentionally skipped source files into a task failure', () => {
expect(
classifyConnectorSyncResult({
Expand Down
5 changes: 4 additions & 1 deletion apps/sim/background/knowledge-connector-sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,7 @@ import type { SyncResult } from '@/connectors/types'

const logger = createLogger('TriggerKnowledgeConnectorSync')

export type ConnectorSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed'
export type ConnectorSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'failed' | 'deferred'

/**
* Separates source-sync failures from expected queue/lock no-ops. Intentional
Expand All @@ -20,6 +20,8 @@ export type ConnectorSyncTaskOutcome = 'completed' | 'partial' | 'skipped' | 'fa
export function classifyConnectorSyncResult(result: SyncResult): ConnectorSyncTaskOutcome {
if (result.skipReason) return 'skipped'
if (result.error) return 'failed'
if (result.deferred && result.docsFailed === 0 && result.processingDispatch.failed === 0)
return 'deferred'
if (result.listingIncomplete || result.docsFailed > 0 || result.processingDispatch.failed > 0)
return 'partial'
return 'completed'
Expand Down Expand Up @@ -61,6 +63,7 @@ export async function executeConnectorSyncJob(payload: unknown) {
logger.info(`[${requestId}] Connector sync completed`, {
connectorId,
outcome: classifyConnectorSyncResult(result),
deferred: result.deferred,
added: result.docsAdded,
updated: result.docsUpdated,
deleted: result.docsDeleted,
Expand Down
Loading
Loading