Skip to content

Commit 5ccf8c1

Browse files
committed
fix(knowledge): unschedule a connector whose credential the source rejected and prompt to reconnect
1 parent 93fdd13 commit 5ccf8c1

13 files changed

Lines changed: 329 additions & 5 deletions

File tree

‎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,7 +4,10 @@ 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'
7+
import {
8+
CREDENTIAL_REMOVED_SYNC_ERROR,
9+
CREDENTIAL_REVOKED_SYNC_ERROR,
10+
} from '@/lib/knowledge/connectors/sync-limits'
811
import { getCanonicalScopesForProvider, getProviderIdFromServiceId } from '@/lib/oauth'
912
import { getMissingRequiredScopes } from '@/lib/oauth/utils'
1013
import { ConnectOAuthModal } from '@/app/workspace/[workspaceId]/components/connect-oauth-modal'
@@ -82,7 +85,10 @@ export function ConnectorRecovery({
8285
const docsUrl = isSearchIndex ? connectorDef?.searchDocsUrl : undefined
8386
const credentialRemoved =
8487
connector.lastSyncError === CREDENTIAL_REMOVED_SYNC_ERROR && !connector.credentialId
85-
const pausedTitle = credentialRemoved
88+
const credentialRevoked =
89+
connector.lastSyncError === CREDENTIAL_REVOKED_SYNC_ERROR && Boolean(connector.credentialId)
90+
const reconnectRequired = credentialRemoved || credentialRevoked
91+
const pausedTitle = reconnectRequired
8692
? 'Reconnect to resume syncing'
8793
: 'Sync paused after repeated failures'
8894

@@ -97,7 +103,7 @@ export function ConnectorRecovery({
97103
}
98104
/>
99105
)}
100-
{connector.status === 'disabled' || credentialRemoved ? (
106+
{connector.status === 'disabled' || reconnectRequired ? (
101107
<SettingsResourceRow
102108
title={
103109
!canEdit

‎apps/sim/connectors/types.ts‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -266,6 +266,7 @@ export const SYNC_SKIP_REASONS = [
266266
'sync_superseded',
267267
'connector_deleted_during_sync',
268268
'credential_missing',
269+
'credential_revoked',
269270
] as const
270271

271272
export type SyncSkipReason = (typeof SYNC_SKIP_REASONS)[number]

‎apps/sim/lib/credentials/draft-hooks.test.ts‎

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -115,6 +115,15 @@ describe('handleReconnectCredential', () => {
115115
expect(mocks.clearDeadFlag).toHaveBeenCalledWith(
116116
getOAuthRefreshCoordinationIdentity('account-new')
117117
)
118+
/** Connectors the rejected credential had unscheduled are due again. */
119+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
120+
expect.objectContaining({
121+
status: 'active',
122+
lastSyncError: null,
123+
consecutiveFailures: 0,
124+
nextSyncAt: new Date('2026-08-14T18:00:00.000Z'),
125+
})
126+
)
118127
expect(auditMockFns.mockRecordAudit).toHaveBeenCalledWith(
119128
expect.objectContaining({
120129
resourceId: 'credential-1',

‎apps/sim/lib/credentials/draft-hooks.ts‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ import { getPostgresConstraintName, getPostgresErrorCode } from '@sim/utils/erro
66
import { generateId } from '@sim/utils/id'
77
import { and, eq, sql } from 'drizzle-orm'
88
import { deleteOrphanedOAuthAccount } from '@/lib/credentials/deletion'
9+
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
910
import { clearOAuthRefreshDeadFlag } from '@/lib/oauth/refresh-coordination'
1011
import { captureServerEvent } from '@/lib/posthog/server'
1112

@@ -75,6 +76,7 @@ export async function handleCreateCredentialFromDraft(params: {
7576
.where(eq(schema.credential.id, existingCredential.id))
7677

7778
await clearOAuthRefreshDeadFlag(accountId)
79+
await resumeConnectorsAfterCredentialReconnect(existingCredential.id, now)
7880

7981
recordAudit({
8082
workspaceId: draft.workspaceId,
@@ -209,6 +211,7 @@ export async function handleReconnectCredential(params: {
209211
)
210212

211213
await clearOAuthRefreshDeadFlag(newAccountId)
214+
await resumeConnectorsAfterCredentialReconnect(draft.credentialId, now)
212215

213216
recordAudit({
214217
workspaceId,

‎apps/sim/lib/credentials/organization-draft.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import { OrchestrationError } from '@/lib/core/orchestration/types'
88
import { resourceScopeCondition } from '@/lib/core/resource-scope.server'
99
import { deleteOrphanedOAuthAccount } from '@/lib/credentials/deletion'
1010
import { getCredentialCreationOrganizationContext } from '@/lib/credentials/organization'
11+
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
1112
import { clearOAuthRefreshDeadFlag } from '@/lib/oauth/refresh-coordination'
1213

1314
/** Completes the exact draft bound to the authenticated provider callback, rechecking current ownership under membership locks. */
@@ -124,6 +125,7 @@ export async function completeOrganizationCredentialDraft(input: {
124125
}
125126
})
126127
await clearOAuthRefreshDeadFlag(input.accountId)
128+
if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(result.credentialId, now)
127129
recordAudit({
128130
actorId: input.userId,
129131
action: result.reconnected
Lines changed: 31 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,31 @@
1+
/**
2+
* @vitest-environment node
3+
*/
4+
import { dbChainMockFns, resetDbChainMock } from '@sim/testing'
5+
import { beforeEach, describe, expect, it, vi } from 'vitest'
6+
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
7+
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
8+
9+
describe('resumeConnectorsAfterCredentialReconnect', () => {
10+
beforeEach(() => {
11+
vi.clearAllMocks()
12+
resetDbChainMock()
13+
})
14+
15+
it('puts only the connectors the rejected credential unscheduled back on schedule', async () => {
16+
const now = new Date('2026-09-22T20:00:00.000Z')
17+
await resumeConnectorsAfterCredentialReconnect('credential-1', now)
18+
expect(dbChainMockFns.set).toHaveBeenCalledWith({
19+
status: 'active',
20+
lastSyncError: null,
21+
consecutiveFailures: 0,
22+
nextSyncAt: now,
23+
updatedAt: now,
24+
})
25+
const guard = JSON.stringify(dbChainMockFns.where.mock.calls.at(-1))
26+
expect(guard).toContain('knowledgeConnector.credentialId')
27+
expect(guard).toContain('credential-1')
28+
expect(guard).toContain('knowledgeConnector.status')
29+
expect(guard).toContain(CREDENTIAL_REVOKED_SYNC_ERROR)
30+
})
31+
})
Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
import { db } from '@sim/db'
2+
import { knowledgeConnector } from '@sim/db/schema'
3+
import { and, eq } from 'drizzle-orm'
4+
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
5+
6+
/**
7+
* Puts the connectors a reconnected credential had unscheduled back on their schedule.
8+
*
9+
* A sync that finds its credential rejected by the source leaves the connector unscheduled
10+
* with {@link CREDENTIAL_REVOKED_SYNC_ERROR}, since retrying cannot help until someone
11+
* authorizes again. Reauthorizing the same credential is that moment: every connector still
12+
* carrying that error is due now, with its failure count cleared. Connectors that were paused
13+
* or disabled for another reason keep their state, and a connector that already moved on is
14+
* left alone.
15+
*/
16+
export async function resumeConnectorsAfterCredentialReconnect(
17+
credentialId: string,
18+
now: Date
19+
): Promise<void> {
20+
await db
21+
.update(knowledgeConnector)
22+
.set({
23+
status: 'active',
24+
lastSyncError: null,
25+
consecutiveFailures: 0,
26+
nextSyncAt: now,
27+
updatedAt: now,
28+
})
29+
.where(
30+
and(
31+
eq(knowledgeConnector.credentialId, credentialId),
32+
eq(knowledgeConnector.status, 'error'),
33+
eq(knowledgeConnector.lastSyncError, CREDENTIAL_REVOKED_SYNC_ERROR)
34+
)
35+
)
36+
}

‎apps/sim/lib/knowledge/connectors/sync-engine.test.ts‎

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@
33
*/
44
import {
55
authOAuthUtilsMock,
6+
authOAuthUtilsMockFns,
67
dbChainMockFns,
78
drizzleOrmMock,
89
flattenMockConditions,
@@ -17,6 +18,7 @@ import { DrizzleQueryError } from 'drizzle-orm/errors'
1718
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
1819
import * as connectorTokens from '@/lib/knowledge/connectors/access-token'
1920
import { executeSync, isConnectorRunnableStatus } from '@/lib/knowledge/connectors/sync-engine'
21+
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
2022
import {
2123
classifySuspectListing,
2224
evaluateListingSafety,
@@ -3007,6 +3009,84 @@ describe('executeSync heartbeats during the listing phase', () => {
30073009
}
30083010
)
30093011

3012+
/** A locked OAuth connector whose token resolution the test controls. */
3013+
function primeOAuthRunUpToToken() {
3014+
const oauthConnector = {
3015+
...CONNECTOR,
3016+
connectorType: 'oauth',
3017+
credentialId: 'cred-1',
3018+
accessMode: 'workspace',
3019+
}
3020+
queueTableRows(schemaMock.knowledgeConnector, [oauthConnector])
3021+
for (let i = 0; i < 20; i++)
3022+
queueTableRows(schemaMock.knowledgeConnector, [
3023+
{ id: 'c-1', connectorArchivedAt: null, connectorDeletedAt: null, kbDeletedAt: null },
3024+
])
3025+
queueTableRows(schemaMock.knowledgeBase, [{ userId: 'u-1', workspaceId: 'ws-1' }])
3026+
dbChainMockFns.returning.mockReset()
3027+
dbChainMockFns.returning.mockResolvedValueOnce([oauthConnector])
3028+
/** The terminal write lands on the row this run still holds. */
3029+
dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'c-1' }])
3030+
const tokenUser = vi
3031+
.spyOn(connectorTokens, 'resolveConnectorTokenUserId')
3032+
.mockResolvedValueOnce('u-1')
3033+
const resolveToken = vi
3034+
.spyOn(connectorTokens, 'resolveConnectorAccessToken')
3035+
.mockResolvedValueOnce(null)
3036+
return () => {
3037+
tokenUser.mockRestore()
3038+
resolveToken.mockRestore()
3039+
}
3040+
}
3041+
3042+
it('unschedules a connector whose credential the source rejected instead of retrying it', async () => {
3043+
const restore = primeOAuthRunUpToToken()
3044+
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValueOnce(
3045+
'invalid_grant'
3046+
)
3047+
try {
3048+
const result = await executeSync('c-1', {
3049+
billingAttribution: { workspaceId: 'ws-1' } as never,
3050+
})
3051+
expect(result.skipReason).toBe('credential_revoked')
3052+
expect(result.error).toBeUndefined()
3053+
expect(authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError).toHaveBeenCalledWith(
3054+
'cred-1'
3055+
)
3056+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
3057+
expect.objectContaining({
3058+
status: 'error',
3059+
nextSyncAt: null,
3060+
lastSyncError: CREDENTIAL_REVOKED_SYNC_ERROR,
3061+
syncLockToken: null,
3062+
syncLockLeaseAt: null,
3063+
})
3064+
)
3065+
expect(dbChainMockFns.set).not.toHaveBeenCalledWith(
3066+
expect.objectContaining({ consecutiveFailures: expect.any(Number) })
3067+
)
3068+
} finally {
3069+
restore()
3070+
}
3071+
})
3072+
3073+
it('keeps the failure ladder for a credential that resolved no token without a terminal error', async () => {
3074+
const restore = primeOAuthRunUpToToken()
3075+
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValueOnce(null)
3076+
try {
3077+
const result = await executeSync('c-1', {
3078+
billingAttribution: { workspaceId: 'ws-1' } as never,
3079+
})
3080+
expect(result.skipReason).toBeUndefined()
3081+
expect(result.error).toContain('Failed to obtain access token')
3082+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
3083+
expect.objectContaining({ status: 'error', consecutiveFailures: 1 })
3084+
)
3085+
} finally {
3086+
restore()
3087+
}
3088+
})
3089+
30103090
it.each([
30113091
{ acl: undefined, incomplete: true },
30123092
{ acl: ['invalid-token'], incomplete: true },

‎apps/sim/lib/knowledge/connectors/sync-engine.ts‎

Lines changed: 66 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,7 @@ import {
5757
CONNECTOR_FAILURE_BACKOFF_CAP_MINUTES,
5858
CONNECTOR_SYNC_MAX_DURATION_SECONDS,
5959
CREDENTIAL_REMOVED_SYNC_ERROR,
60+
CREDENTIAL_REVOKED_SYNC_ERROR,
6061
connectorFailureBackoffMinutes,
6162
MAX_CONSECUTIVE_FAILURES,
6263
} from '@/lib/knowledge/connectors/sync-limits'
@@ -88,6 +89,7 @@ import {
8889
import { hardDeleteDocuments } from '@/lib/knowledge/documents/service'
8990
import { getRetryAfterMs, isRateLimitError } from '@/lib/knowledge/documents/utils'
9091
import { ensureSourceVectorIndex } from '@/lib/knowledge/search/source-vector-indexes'
92+
import { getCredentialTerminalRefreshError } from '@/lib/oauth/credential-service'
9193
import { connectorHasAuthSource } from '@/connectors/auth'
9294
import { CONNECTOR_REGISTRY } from '@/connectors/registry.server'
9395
import type {
@@ -743,9 +745,26 @@ export function buildSyncSuccessUpdate(
743745
}
744746
}
745747

748+
/**
749+
* A credential the source rejected outright: the refresh path recorded a terminal error for
750+
* it, so no retry can produce a token until the credential is reauthorized.
751+
*/
752+
export class ConnectorCredentialRevokedError extends Error {
753+
constructor(
754+
readonly credentialId: string,
755+
readonly errorCode: string
756+
) {
757+
super(`Credential ${credentialId} was rejected by the source (${errorCode})`)
758+
this.name = 'ConnectorCredentialRevokedError'
759+
}
760+
}
761+
746762
/**
747763
* Resolves the token a connector syncs with, failing loudly where the shared
748764
* resolver reports "no token" — a sync has no reconnect prompt to fall back to.
765+
* A credential the source has rejected outright fails as
766+
* {@link ConnectorCredentialRevokedError}, so the run can unschedule the
767+
* connector instead of walking the failure ladder toward a retry that cannot help.
749768
*/
750769
async function resolveAccessToken(
751770
connector: { credentialId: string | null; encryptedApiKey: string | null },
@@ -770,6 +789,13 @@ async function resolveAccessToken(
770789
userId,
771790
authMode: connectorConfig.auth.mode,
772791
})
792+
const terminalError =
793+
connectorConfig.auth.mode === 'oauth' && connector.credentialId
794+
? await getCredentialTerminalRefreshError(connector.credentialId)
795+
: null
796+
if (terminalError && connector.credentialId) {
797+
throw new ConnectorCredentialRevokedError(connector.credentialId, terminalError)
798+
}
773799
throw new Error(`Failed to obtain access token for credential ${connector.credentialId}`)
774800
}
775801

@@ -1427,6 +1453,46 @@ export async function executeSync(
14271453
}
14281454
}
14291455

1456+
if (error instanceof ConnectorCredentialRevokedError) {
1457+
/**
1458+
* Retrying cannot help until the credential is reauthorized, so the
1459+
* connector leaves its schedule with a reconnect prompt instead of
1460+
* climbing the failure ladder toward the same rejection. Reauthorizing
1461+
* the credential puts it back on schedule. The run itself is a skip:
1462+
* nothing about the source failed, and a sync that cannot start is not
1463+
* an incident to page on.
1464+
*/
1465+
logger.warn('Sync unscheduled: the source rejected the connector credential', {
1466+
connectorId,
1467+
credentialId: error.credentialId,
1468+
errorCode: error.errorCode,
1469+
})
1470+
try {
1471+
await completeSyncLog(syncLogId, 'failed', result, {
1472+
errorMessage: CREDENTIAL_REVOKED_SYNC_ERROR,
1473+
})
1474+
const landed = await writeTerminalConnectorState(
1475+
connectorId,
1476+
syncLogId,
1477+
buildSyncUnscheduledUpdate(new Date(), CREDENTIAL_REVOKED_SYNC_ERROR)
1478+
)
1479+
if (!landed) {
1480+
logger.warn(
1481+
'Unschedule discarded — connector was reclaimed while this run was executing',
1482+
{ connectorId, syncLogId }
1483+
)
1484+
}
1485+
} catch (recoveryError) {
1486+
logger.error('Failed to unschedule the connector', {
1487+
connectorId,
1488+
error:
1489+
getConnectorFailureDiagnostic(recoveryError)?.message ??
1490+
toError(recoveryError).message,
1491+
})
1492+
}
1493+
return { ...result, skipReason: 'credential_revoked' }
1494+
}
1495+
14301496
const diagnostic = getConnectorFailureDiagnostic(error)
14311497
const errorMessage = diagnostic?.message ?? toError(error).message
14321498
const retryAfterMs = getRetryAfterMs(error)

‎apps/sim/lib/knowledge/connectors/sync-limits.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,13 @@ export const MAX_CONSECUTIVE_FAILURES = 10
2525
export const CREDENTIAL_REMOVED_SYNC_ERROR =
2626
'Credential removed. Reconnect the connector to resume syncing.'
2727

28+
/**
29+
* The error a connector carries once the source rejects its credential outright (a revoked or
30+
* expired grant, not a passing failure); cleared by reauthorizing that credential.
31+
*/
32+
export const CREDENTIAL_REVOKED_SYNC_ERROR =
33+
'The source no longer accepts this credential. Reconnect it to resume syncing.'
34+
2835
/**
2936
* The error a connector carries once {@link MAX_CONSECUTIVE_FAILURES} disables it.
3037
*

0 commit comments

Comments
 (0)