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
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,10 @@ import { useEffect, useState } from 'react'
import { Chip } from '@sim/emcn'
import type { ConnectorData } from '@/lib/api/contracts/knowledge/connectors'
import { type ResourceScope, resourceScopeFields } from '@/lib/core/resource-scope'
import { CREDENTIAL_REMOVED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
import {
CREDENTIAL_REMOVED_SYNC_ERROR,
CREDENTIAL_REVOKED_SYNC_ERROR,
} from '@/lib/knowledge/connectors/sync-limits'
import { getCanonicalScopesForProvider, getProviderIdFromServiceId } from '@/lib/oauth'
import { getMissingRequiredScopes } from '@/lib/oauth/utils'
import { ConnectOAuthModal } from '@/app/workspace/[workspaceId]/components/connect-oauth-modal'
Expand Down Expand Up @@ -82,7 +85,10 @@ export function ConnectorRecovery({
const docsUrl = isSearchIndex ? connectorDef?.searchDocsUrl : undefined
const credentialRemoved =
connector.lastSyncError === CREDENTIAL_REMOVED_SYNC_ERROR && !connector.credentialId
const pausedTitle = credentialRemoved
const credentialRevoked =
connector.lastSyncError === CREDENTIAL_REVOKED_SYNC_ERROR && Boolean(connector.credentialId)
const reconnectRequired = credentialRemoved || credentialRevoked
const pausedTitle = reconnectRequired
? 'Reconnect to resume syncing'
: 'Sync paused after repeated failures'

Expand All @@ -97,7 +103,7 @@ export function ConnectorRecovery({
}
/>
)}
{connector.status === 'disabled' || credentialRemoved ? (
{connector.status === 'disabled' || reconnectRequired ? (
<SettingsResourceRow
title={
!canEdit
Expand Down
1 change: 1 addition & 0 deletions apps/sim/connectors/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -266,6 +266,7 @@ export const SYNC_SKIP_REASONS = [
'sync_superseded',
'connector_deleted_during_sync',
'credential_missing',
'credential_revoked',
] as const

export type SyncSkipReason = (typeof SYNC_SKIP_REASONS)[number]
Expand Down
10 changes: 10 additions & 0 deletions apps/sim/lib/credentials/draft-hooks.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,7 @@ describe('handleReconnectCredential', () => {
{ id: 'credential-1', accountId: null, displayName: 'Renamed Gmail' },
])
queueTableRows(schemaMock.credential, [])
queueTableRows(schemaMock.account, [{ providerId: 'gmail', accountId: 'subject-new' }])

await handleReconnectCredential({
draft: { credentialId: 'credential-1' },
Expand All @@ -115,6 +116,15 @@ describe('handleReconnectCredential', () => {
expect(mocks.clearDeadFlag).toHaveBeenCalledWith(
getOAuthRefreshCoordinationIdentity('account-new')
)
/** Connectors the rejected credential had unscheduled are due again. */
expect(dbChainMockFns.set).toHaveBeenCalledWith(
expect.objectContaining({
status: 'active',
lastSyncError: null,
consecutiveFailures: 0,
nextSyncAt: new Date('2026-08-14T18:00:00.000Z'),
})
)
expect(auditMockFns.mockRecordAudit).toHaveBeenCalledWith(
expect.objectContaining({
resourceId: 'credential-1',
Expand Down
3 changes: 3 additions & 0 deletions apps/sim/lib/credentials/draft-hooks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import { getPostgresConstraintName, getPostgresErrorCode } from '@sim/utils/erro
import { generateId } from '@sim/utils/id'
import { and, eq, sql } from 'drizzle-orm'
import { deleteOrphanedOAuthAccount } from '@/lib/credentials/deletion'
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
import { clearOAuthRefreshDeadFlag } from '@/lib/oauth/refresh-coordination'
import { captureServerEvent } from '@/lib/posthog/server'

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

await clearOAuthRefreshDeadFlag(accountId)
await resumeConnectorsAfterCredentialReconnect(accountId, now)

recordAudit({
workspaceId: draft.workspaceId,
Expand Down Expand Up @@ -209,6 +211,7 @@ export async function handleReconnectCredential(params: {
)

await clearOAuthRefreshDeadFlag(newAccountId)
await resumeConnectorsAfterCredentialReconnect(newAccountId, now)

recordAudit({
workspaceId,
Expand Down
2 changes: 2 additions & 0 deletions apps/sim/lib/credentials/organization-draft.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { OrchestrationError } from '@/lib/core/orchestration/types'
import { resourceScopeCondition } from '@/lib/core/resource-scope.server'
import { deleteOrphanedOAuthAccount } from '@/lib/credentials/deletion'
import { getCredentialCreationOrganizationContext } from '@/lib/credentials/organization'
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
import { clearOAuthRefreshDeadFlag } from '@/lib/oauth/refresh-coordination'

/** Completes the exact draft bound to the authenticated provider callback, rechecking current ownership under membership locks. */
Expand Down Expand Up @@ -124,6 +125,7 @@ export async function completeOrganizationCredentialDraft(input: {
}
})
await clearOAuthRefreshDeadFlag(input.accountId)
if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(input.accountId, now)
recordAudit({
actorId: input.userId,
action: result.reconnected
Expand Down
52 changes: 52 additions & 0 deletions apps/sim/lib/knowledge/connectors/credential-recovery.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
/**
* @vitest-environment node
*/
import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing'
import { beforeEach, describe, expect, it, vi } from 'vitest'
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'

const RESUMED = {
status: 'active',
lastSyncError: null,
consecutiveFailures: 0,
}

describe('resumeConnectorsAfterCredentialReconnect', () => {
const now = new Date('2026-09-22T20:00:00.000Z')

beforeEach(() => {
vi.clearAllMocks()
resetDbChainMock()
})

it('resumes the connectors of every credential on the reconnected account', async () => {
queueTableRows(schemaMock.account, [
{ providerId: 'confluence', providerAccountId: 'subject-1' },
])
await resumeConnectorsAfterCredentialReconnect('account-1', now)
expect(dbChainMockFns.set).toHaveBeenCalledWith({ ...RESUMED, nextSyncAt: now, updatedAt: now })
const guard = JSON.stringify(dbChainMockFns.where.mock.calls.at(-1))
expect(guard).toContain('knowledgeConnector.credentialId')
expect(guard).toContain(CREDENTIAL_REVOKED_SYNC_ERROR)
const credentials = JSON.stringify(dbChainMockFns.where.mock.calls)
expect(credentials).toContain('"left":"credential.accountId","right":"account-1"')
})

it('resumes across the Slack installation, whose sibling accounts share the repaired chain', async () => {
queueTableRows(schemaMock.account, [
{ providerId: 'slack', providerAccountId: 'TEXAMPLE-usr_U1' },
])
await resumeConnectorsAfterCredentialReconnect('account-1', now)
expect(dbChainMockFns.set).toHaveBeenCalledWith({ ...RESUMED, nextSyncAt: now, updatedAt: now })
const conditions = JSON.stringify(dbChainMockFns.where.mock.calls)
expect(conditions).toContain('"pattern":"TEXAMPLE-%"')
expect(conditions).not.toContain('"left":"credential.accountId","right":"account-1"')
})

it('does nothing for an account that no longer exists', async () => {
queueTableRows(schemaMock.account, [])
await resumeConnectorsAfterCredentialReconnect('account-gone', now)
expect(dbChainMockFns.update).not.toHaveBeenCalled()
})
})
57 changes: 57 additions & 0 deletions apps/sim/lib/knowledge/connectors/credential-recovery.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
import { db } from '@sim/db'
import { account, credential, knowledgeConnector } from '@sim/db/schema'
import { and, eq, inArray } from 'drizzle-orm'
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
import { extractSlackTeamId, installationFilter, isSlackProvider } from '@/lib/oauth/slack'

/**
* Puts the connectors a reconnected account had unscheduled back on their schedule.
*
* A sync that finds its credential rejected by the source leaves the connector unscheduled
* with {@link CREDENTIAL_REVOKED_SYNC_ERROR}, since retrying cannot help until someone
* authorizes again. Reauthorizing the account is that moment: every connector on a credential
* of that account still carrying the error is due now, with its failure count cleared. A Slack
* reauthorization repairs the installation's shared token chain, so the connectors on every
* credential of the installation's sibling accounts are due as well. Connectors paused or
* disabled for another reason keep their state, and a connector that already moved on is left
* alone.
*/
export async function resumeConnectorsAfterCredentialReconnect(
accountId: string,
now: Date
): Promise<void> {
const [reconnected] = await db
.select({ providerId: account.providerId, providerAccountId: account.accountId })
.from(account)
.where(eq(account.id, accountId))
.limit(1)
if (!reconnected) return
const slackTeamId = isSlackProvider(reconnected.providerId)
? extractSlackTeamId(reconnected.providerAccountId)
: null
const repairedAccounts = slackTeamId
? inArray(
credential.accountId,
db.select({ id: account.id }).from(account).where(installationFilter(slackTeamId))
)
: eq(credential.accountId, accountId)
await db
.update(knowledgeConnector)
.set({
status: 'active',
lastSyncError: null,
consecutiveFailures: 0,
nextSyncAt: now,
updatedAt: now,
})
.where(
and(
inArray(
knowledgeConnector.credentialId,
db.select({ id: credential.id }).from(credential).where(repairedAccounts)
),
eq(knowledgeConnector.status, 'error'),
Comment thread
waleedlatif1 marked this conversation as resolved.
eq(knowledgeConnector.lastSyncError, CREDENTIAL_REVOKED_SYNC_ERROR)
)
)
}
127 changes: 127 additions & 0 deletions apps/sim/lib/knowledge/connectors/sync-engine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
*/
import {
authOAuthUtilsMock,
authOAuthUtilsMockFns,
dbChainMockFns,
drizzleOrmMock,
flattenMockConditions,
Expand All @@ -17,6 +18,7 @@ import { DrizzleQueryError } from 'drizzle-orm/errors'
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
import * as connectorTokens from '@/lib/knowledge/connectors/access-token'
import { executeSync, isConnectorRunnableStatus } from '@/lib/knowledge/connectors/sync-engine'
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
import {
classifySuspectListing,
evaluateListingSafety,
Expand Down Expand Up @@ -3007,6 +3009,131 @@ describe('executeSync heartbeats during the listing phase', () => {
}
)

/** A locked OAuth connector whose token resolution the test controls. */
function primeOAuthRunUpToToken() {
const oauthConnector = {
...CONNECTOR,
connectorType: 'oauth',
credentialId: 'cred-1',
accessMode: 'workspace',
}
queueTableRows(schemaMock.knowledgeConnector, [oauthConnector])
for (let i = 0; i < 20; i++)
queueTableRows(schemaMock.knowledgeConnector, [
{ id: 'c-1', connectorArchivedAt: null, connectorDeletedAt: null, kbDeletedAt: null },
])
queueTableRows(schemaMock.knowledgeBase, [{ userId: 'u-1', workspaceId: 'ws-1' }])
dbChainMockFns.returning.mockReset()
dbChainMockFns.returning.mockResolvedValueOnce([oauthConnector])
/** The terminal write lands on the row this run still holds. */
dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'c-1' }])
const tokenUser = vi
.spyOn(connectorTokens, 'resolveConnectorTokenUserId')
.mockResolvedValueOnce('u-1')
const resolveToken = vi
.spyOn(connectorTokens, 'resolveConnectorAccessToken')
.mockResolvedValueOnce(null)
return () => {
tokenUser.mockRestore()
resolveToken.mockRestore()
}
}

it('unschedules a connector whose credential the source rejected instead of retrying it', async () => {
const restore = primeOAuthRunUpToToken()
/** Rejected at token resolution and still rejected when the run records its outcome. */
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce('invalid_grant')
try {
const result = await executeSync('c-1', {
billingAttribution: { workspaceId: 'ws-1' } as never,
})
expect(result.skipReason).toBe('credential_revoked')
expect(result.error).toBeUndefined()
expect(authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError).toHaveBeenCalledWith(
'cred-1'
)
expect(dbChainMockFns.set).toHaveBeenCalledWith(
expect.objectContaining({
status: 'error',
nextSyncAt: null,
lastSyncError: CREDENTIAL_REVOKED_SYNC_ERROR,
syncLockToken: null,
syncLockLeaseAt: null,
})
)
expect(dbChainMockFns.set).not.toHaveBeenCalledWith(
expect.objectContaining({ consecutiveFailures: expect.any(Number) })
)
} finally {
restore()
}
})

it('takes the failure ladder when the credential was reauthorized while the run was failing', async () => {
const restore = primeOAuthRunUpToToken()
/** Rejected at token resolution, repaired by the time the run records its outcome. */
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce(null)
try {
const result = await executeSync('c-1', {
billingAttribution: { workspaceId: 'ws-1' } as never,
})
expect(result.skipReason).toBeUndefined()
expect(result.error).toContain('rejected by the source')
expect(dbChainMockFns.set).toHaveBeenCalledWith(
expect.objectContaining({ status: 'error', consecutiveFailures: 1 })
)
expect(dbChainMockFns.set).not.toHaveBeenCalledWith(
expect.objectContaining({ lastSyncError: CREDENTIAL_REVOKED_SYNC_ERROR })
)
} finally {
restore()
}
})

it('reports a run that could not record the unschedule as failed, not skipped', async () => {
const restore = primeOAuthRunUpToToken()
/** Rejected at token resolution and still rejected when the run records its outcome. */
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce('invalid_grant')
/** The terminal write fails after the lock CAS consumed the first result. */
dbChainMockFns.returning.mockReset()
dbChainMockFns.returning.mockResolvedValueOnce([
{ ...CONNECTOR, connectorType: 'oauth', credentialId: 'cred-1', accessMode: 'workspace' },
])
dbChainMockFns.returning.mockRejectedValueOnce(new Error('connection reset'))
try {
const result = await executeSync('c-1', {
billingAttribution: { workspaceId: 'ws-1' } as never,
})
expect(result.skipReason).toBeUndefined()
expect(result.error).toContain('connection reset')
} finally {
restore()
}
})

it('keeps the failure ladder for a credential that resolved no token without a terminal error', async () => {
const restore = primeOAuthRunUpToToken()
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValueOnce(null)
try {
const result = await executeSync('c-1', {
billingAttribution: { workspaceId: 'ws-1' } as never,
})
expect(result.skipReason).toBeUndefined()
expect(result.error).toContain('Failed to obtain access token')
expect(dbChainMockFns.set).toHaveBeenCalledWith(
expect.objectContaining({ status: 'error', consecutiveFailures: 1 })
)
} finally {
restore()
}
})

it.each([
{ acl: undefined, incomplete: true },
{ acl: ['invalid-token'], incomplete: true },
Expand Down
Loading
Loading