Skip to content

Commit 745cdbe

Browse files
authored
fix(knowledge): make connector syncs resilient to removed credentials, partial outcomes and slow completion counts (#8126)
* fix(knowledge): leave a connector whose credential was removed unscheduled instead of walking the failure ladder to disabled * fix(knowledge): report a partial connector sync as an outcome instead of failing the run * fix(knowledge): count a connector's documents before the completion locks and index the live set * fix(billing): serve execution admission from the usage gate cache instead of re-summing the billing period per run * fix(knowledge): share the unscheduled connector write, scope credential removal to content-engine modes, and keep both sync holds visible * fix(knowledge): keep the corpus size on a failed completion count, unschedule only connectors left without an auth source, and give fixture connectors a credential * fix(knowledge): skip a required-key connector without a key and write the unscheduled state only over the row the run observed * fix(knowledge): release a dispatch-marked pending row when its credential is missing * fix(knowledge): assert the billing owner before the credential-missing terminal write * fix(knowledge): stop offering reconnect once a connector holds a credential again
1 parent b37ce85 commit 745cdbe

44 files changed

Lines changed: 29030 additions & 140 deletions

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/v1/knowledge/search/route.test.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,9 @@ vi.mock('@/lib/billing/calculations/usage-monitor', () => ({
7070
vi.mock('@/lib/billing/core/billing-attribution', () => ({
7171
resolveBillingAttribution: mockResolveBillingAttribution,
7272
resolveSystemBillingAttribution: mockResolveSystemBillingAttribution,
73-
checkAttributedUsageLimits: vi.fn().mockResolvedValue({ isExceeded: false }),
73+
}))
74+
vi.mock('@/lib/billing/core/usage-gate-cache', () => ({
75+
checkSearchUsageLimits: vi.fn().mockResolvedValue({ isExceeded: false }),
7476
}))
7577

7678
vi.mock('@/lib/knowledge/embeddings', () => ({

‎apps/sim/app/api/v1/knowledge/search/route.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,10 +2,10 @@ import { type NextRequest, NextResponse } from 'next/server'
22
import { v1KnowledgeSearchContract } from '@/lib/api/contracts/v1/knowledge'
33
import { parseRequest } from '@/lib/api/server'
44
import {
5-
checkAttributedUsageLimits,
65
resolveBillingAttribution,
76
resolveSystemBillingAttribution,
87
} from '@/lib/billing/core/billing-attribution'
8+
import { checkSearchUsageLimits } from '@/lib/billing/core/usage-gate-cache'
99
import { withRouteHandler } from '@/lib/core/utils/with-route-handler'
1010
import { ALL_TAG_SLOTS } from '@/lib/knowledge/constants'
1111
import { toKbEmbeddingDimensions } from '@/lib/knowledge/embedding-models'
@@ -82,7 +82,7 @@ export const POST = withRouteHandler(async (request: NextRequest) => {
8282
* keys resolve their system actor and immutable payer from one workspace read.
8383
*/
8484
if (billingAttribution) {
85-
const usage = await checkAttributedUsageLimits(billingAttribution)
85+
const usage = await checkSearchUsageLimits(billingAttribution)
8686
if (usage.isExceeded) {
8787
return NextResponse.json(
8888
{ error: usage.message || 'Usage limit exceeded. Please upgrade your plan to continue.' },

‎apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connector-recovery.tsx‎

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import { useEffect, useState } from 'react'
44
import { Chip } from '@sim/emcn'
55
import type { ConnectorData } from '@/lib/api/contracts/knowledge/connectors'
66
import { type ResourceScope, resourceScopeFields } from '@/lib/core/resource-scope'
7+
import { CREDENTIAL_REMOVED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
78
import { getCanonicalScopesForProvider, getProviderIdFromServiceId } from '@/lib/oauth'
89
import { getMissingRequiredScopes } from '@/lib/oauth/utils'
910
import { ConnectOAuthModal } from '@/app/workspace/[workspaceId]/components/connect-oauth-modal'
@@ -79,6 +80,11 @@ export function ConnectorRecovery({
7980
}
8081

8182
const docsUrl = isSearchIndex ? connectorDef?.searchDocsUrl : undefined
83+
const credentialRemoved =
84+
connector.lastSyncError === CREDENTIAL_REMOVED_SYNC_ERROR && !connector.credentialId
85+
const pausedTitle = credentialRemoved
86+
? 'Reconnect to resume syncing'
87+
: 'Sync paused after repeated failures'
8288

8389
return (
8490
<>
@@ -91,16 +97,16 @@ export function ConnectorRecovery({
9197
}
9298
/>
9399
)}
94-
{connector.status === 'disabled' ? (
100+
{connector.status === 'disabled' || credentialRemoved ? (
95101
<SettingsResourceRow
96102
title={
97103
!canEdit
98-
? 'Sync paused after repeated failures'
104+
? pausedTitle
99105
: requiresAccountSettings
100106
? 'Update the source account, then resume syncing'
101107
: serviceId
102108
? 'Reconnect to resume syncing'
103-
: 'Sync paused after repeated failures'
109+
: pausedTitle
104110
}
105111
trailing={
106112
canEdit && requiresAccountSettings && onEdit ? (

‎apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connectors-section.test.tsx‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { afterEach, describe, expect, it, vi } from 'vitest'
88
import type { SyncLogData } from '@/lib/api/contracts/knowledge/connectors'
99
import {
1010
CONNECTOR_SYNC_STALE_LOCK_TTL_MS,
11+
CREDENTIAL_REMOVED_SYNC_ERROR,
1112
MEMBER_SYNC_STALE_LOCK_TTL_MS,
1213
} from '@/lib/knowledge/connectors/sync-limits'
1314

@@ -480,6 +481,21 @@ describe('Connector credential reauthorization', () => {
480481
expect(container.textContent).toContain('No connected sources yet.')
481482
})
482483

484+
it('offers reconnect for a connector whose credential was removed', () => {
485+
oauthCredentialsState.current = [
486+
{ id: 'credential-1', name: 'Workspace Slack', provider: 'slack-custom' },
487+
]
488+
const container = renderSection(
489+
makeConnector({
490+
status: 'error',
491+
credentialId: null,
492+
lastSyncError: CREDENTIAL_REMOVED_SYNC_ERROR,
493+
})
494+
)
495+
expect(container.textContent).toContain('Reconnect to resume syncing')
496+
expect(findButton(container, 'Reconnect')).not.toBeDisabled()
497+
})
498+
483499
it('reauthorizes with the resolved credential provider and identity', () => {
484500
oauthCredentialsState.current = [
485501
{ id: 'credential-1', name: 'Workspace Slack', provider: 'slack-custom' },

‎apps/sim/background/drain-governed-subject.test.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,9 +49,11 @@ vi.mock('@/enrichments/run', () => ({
4949
}))
5050
vi.mock('@/lib/billing/core/billing-attribution', () => ({
5151
assertBillingAttributionSnapshot: (snapshot: unknown) => snapshot,
52-
checkAttributedUsageLimits: mocks.checkAttributedUsageLimits,
5352
toBillingContext: () => ({}),
5453
}))
54+
vi.mock('@/lib/billing/core/usage-gate-cache', () => ({
55+
checkExecutionUsageLimits: mocks.checkAttributedUsageLimits,
56+
}))
5557
vi.mock('@/lib/table/rows/secret-provenance', () => ({
5658
createExactEmptyTableRowSecretProvenance: () => ({ complete: true, columns: {} }),
5759
createTableRowSecretProvenanceFromRegistry: () => ({ complete: true, columns: {} }),

‎apps/sim/background/enrichment-capability-subject.test.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -58,9 +58,11 @@ vi.mock('@/enrichments/run', () => ({
5858
}))
5959
vi.mock('@/lib/billing/core/billing-attribution', () => ({
6060
assertBillingAttributionSnapshot: vi.fn((value) => value),
61-
checkAttributedUsageLimits: mocks.checkAttributedUsageLimits,
6261
toBillingContext: vi.fn(() => ({})),
6362
}))
63+
vi.mock('@/lib/billing/core/usage-gate-cache', () => ({
64+
checkExecutionUsageLimits: mocks.checkAttributedUsageLimits,
65+
}))
6466
vi.mock('@/lib/table/rows/secret-provenance', () => ({
6567
createExactEmptyTableRowSecretProvenance: vi.fn(() => undefined),
6668
createTableRowSecretProvenanceFromRegistry: vi.fn(() => undefined),

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

Lines changed: 17 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,12 +5,16 @@
55
import { AbortTaskRunError } from '@trigger.dev/sdk'
66
import { beforeEach, describe, expect, it, vi } from 'vitest'
77

8-
const { mockAssertConnectorSyncPayload, mockExecuteSync, mockTask } = vi.hoisted(() => ({
8+
const { mockAssertConnectorSyncPayload, mockExecuteSync, mockTask, mockWarn } = vi.hoisted(() => ({
9+
mockWarn: vi.fn(),
910
mockAssertConnectorSyncPayload: vi.fn(),
1011
mockExecuteSync: vi.fn(),
1112
mockTask: vi.fn((config) => config),
1213
}))
1314

15+
vi.mock('@sim/logger', () => ({
16+
createLogger: () => ({ info: vi.fn(), warn: mockWarn, error: vi.fn(), debug: vi.fn() }),
17+
}))
1418
vi.mock('@trigger.dev/sdk', () => ({
1519
task: mockTask,
1620
AbortTaskRunError: class AbortTaskRunError extends Error {},
@@ -114,7 +118,7 @@ describe('knowledge connector sync worker', () => {
114118
})
115119
})
116120

117-
it('fails visibly without retrying an already-persisted partial sync', async () => {
121+
it('returns a partial sync as an outcome instead of failing the run', async () => {
118122
mockAssertConnectorSyncPayload.mockReturnValue({
119123
connectorId: 'connector-1',
120124
requestId: 'request-1',
@@ -136,8 +140,15 @@ describe('knowledge connector sync worker', () => {
136140
billingAttribution: BILLING_ATTRIBUTION,
137141
})
138142

139-
await expect(run).rejects.toBeInstanceOf(AbortTaskRunError)
140-
await expect(run).rejects.toThrow('Connector sync partially failed')
143+
await expect(run).resolves.toMatchObject({
144+
success: false,
145+
outcome: 'partial',
146+
docsFailed: 1,
147+
processingDispatch: { failed: 1 },
148+
})
149+
expect(mockWarn).toHaveBeenCalledWith(
150+
expect.stringContaining('1 source failures, 1 dispatch failures')
151+
)
141152
})
142153

143154
it('completes a durably scheduled capacity wait while preserving existing source failures', async () => {
@@ -163,7 +174,7 @@ describe('knowledge connector sync worker', () => {
163174
deferred: waiting.deferred,
164175
})
165176
mockExecuteSync.mockResolvedValue({ ...waiting, docsFailed: 1 })
166-
await expect(executeConnectorSyncJob({})).rejects.toThrow('partially failed')
177+
expect(await executeConnectorSyncJob({})).toMatchObject({ outcome: 'partial', success: false })
167178
mockExecuteSync.mockResolvedValue({ ...waiting, error: 'Retry persistence failed' })
168179
await expect(executeConnectorSyncJob({})).rejects.toThrow('Retry persistence failed')
169180
})
@@ -266,6 +277,7 @@ describe('knowledge connector sync worker', () => {
266277
outcome: 'partial',
267278
listingIncomplete: true,
268279
})
280+
expect(mockWarn).not.toHaveBeenCalled()
269281
})
270282

271283
it('classifies a persisted connector error as a failed task', () => {

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

Lines changed: 14 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -27,17 +27,6 @@ export function classifyConnectorSyncResult(result: SyncResult): ConnectorSyncTa
2727
return 'completed'
2828
}
2929

30-
function formatConnectorSyncFailure(
31-
connectorId: string,
32-
result: SyncResult,
33-
outcome: Extract<ConnectorSyncTaskOutcome, 'partial' | 'failed'>
34-
): string {
35-
if (outcome === 'failed') {
36-
return `Connector sync failed for ${connectorId}: ${result.error}`
37-
}
38-
return `Connector sync partially failed for ${connectorId}: ${result.docsFailed} source failures, ${result.processingDispatch.failed} dispatch failures`
39-
}
40-
4130
export async function executeConnectorSyncJob(payload: unknown) {
4231
const {
4332
connectorId,
@@ -60,9 +49,10 @@ export async function executeConnectorSyncJob(payload: unknown) {
6049
dispatchToken,
6150
})
6251

52+
const outcome = classifyConnectorSyncResult(result)
6353
logger.info(`[${requestId}] Connector sync completed`, {
6454
connectorId,
65-
outcome: classifyConnectorSyncResult(result),
55+
outcome,
6656
deferred: result.deferred,
6757
added: result.docsAdded,
6858
updated: result.docsUpdated,
@@ -75,17 +65,20 @@ export async function executeConnectorSyncJob(payload: unknown) {
7565
processingDispatchFailed: result.processingDispatch.failed,
7666
})
7767

78-
const outcome = classifyConnectorSyncResult(result)
79-
if (outcome === 'failed' || result.docsFailed > 0 || result.processingDispatch.failed > 0) {
68+
if (outcome === 'failed') {
8069
/**
81-
* `executeSync` has already persisted its terminal state. Source failures
82-
* preserve the previous incremental watermark so the next connector pass
83-
* replays them; dispatch failures remain eligible for the stuck-document
84-
* sweep. Retrying this whole task immediately would duplicate a large
85-
* fan-out, so fail visibly without retrying the completed transaction.
70+
* `executeSync` has already persisted its terminal state, and retrying this
71+
* whole task would duplicate a large fan-out, so fail visibly without
72+
* retrying the completed transaction. A partial sync is not a failed run:
73+
* its source failures keep the previous incremental watermark so the next
74+
* connector pass replays them, its dispatch failures stay eligible for the
75+
* stuck-document sweep, and the outcome rides on the return value.
8676
*/
87-
throw new AbortTaskRunError(
88-
formatConnectorSyncFailure(connectorId, result, outcome === 'failed' ? 'failed' : 'partial')
77+
throw new AbortTaskRunError(`Connector sync failed for ${connectorId}: ${result.error}`)
78+
}
79+
if (outcome === 'partial' && (result.docsFailed > 0 || result.processingDispatch.failed > 0)) {
80+
logger.warn(
81+
`[${requestId}] Connector sync partially failed for ${connectorId}: ${result.docsFailed} source failures, ${result.processingDispatch.failed} dispatch failures`
8982
)
9083
}
9184

‎apps/sim/background/workflow-column-execution.ts‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,9 +10,9 @@ import { and, eq, isNull, or } from 'drizzle-orm'
1010
import {
1111
assertBillingAttributionSnapshot,
1212
type BillingAttributionSnapshot,
13-
checkAttributedUsageLimits,
1413
toBillingContext,
1514
} from '@/lib/billing/core/billing-attribution'
15+
import { checkExecutionUsageLimits } from '@/lib/billing/core/usage-gate-cache'
1616
import { checkAndBillPayerOverageThreshold } from '@/lib/billing/threshold-billing'
1717
import { isRetryableInfrastructureError } from '@/lib/core/errors/retryable-infrastructure'
1818
import {
@@ -575,7 +575,7 @@ async function runWorkflowAndWriteTerminal(
575575
* Gate the exact workspace payer and member cap before hosted-key cost.
576576
* A denial clears the cell pre-stamp and surfaces the upgrade state.
577577
*/
578-
const usage = await checkAttributedUsageLimits(enrichmentBillingAttribution)
578+
const usage = await checkExecutionUsageLimits(enrichmentBillingAttribution)
579579
if (usage.isExceeded) {
580580
logger.warn(
581581
`Usage limit reached — halting enrichment (table=${tableId} row=${rowId} group=${groupId})`

‎apps/sim/background/workflow-group-governed-subject.test.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,9 +93,11 @@ vi.mock('@/executor/utils/resolved-secret-trace-registry', () => ({
9393
}))
9494
vi.mock('@/lib/billing/core/billing-attribution', () => ({
9595
assertBillingAttributionSnapshot: (snapshot: unknown) => snapshot,
96-
checkAttributedUsageLimits: async () => ({ isExceeded: false }),
9796
toBillingContext: () => ({}),
9897
}))
98+
vi.mock('@/lib/billing/core/usage-gate-cache', () => ({
99+
checkExecutionUsageLimits: async () => ({ isExceeded: false }),
100+
}))
99101
/** Real pacing would sleep jittered backoff against the global db mock. */
100102
vi.mock('@/lib/core/rate-limiter/rate-limiter', () => ({
101103
RateLimiter: class {

0 commit comments

Comments
 (0)