From 5ccf8c1bb664fb506bd4634c70c8e29131c08d6b Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 22 Sep 2026 13:43:16 -0700 Subject: [PATCH 1/2] fix(knowledge): unschedule a connector whose credential the source rejected and prompt to reconnect --- .../connectors-section/connector-recovery.tsx | 12 ++- apps/sim/connectors/types.ts | 1 + apps/sim/lib/credentials/draft-hooks.test.ts | 9 +++ apps/sim/lib/credentials/draft-hooks.ts | 3 + .../sim/lib/credentials/organization-draft.ts | 2 + .../connectors/credential-recovery.test.ts | 31 +++++++ .../connectors/credential-recovery.ts | 36 +++++++++ .../knowledge/connectors/sync-engine.test.ts | 80 +++++++++++++++++++ .../lib/knowledge/connectors/sync-engine.ts | 66 +++++++++++++++ .../lib/knowledge/connectors/sync-limits.ts | 7 ++ apps/sim/lib/oauth/credential-service.test.ts | 43 ++++++++++ apps/sim/lib/oauth/credential-service.ts | 42 +++++++++- .../src/mocks/auth-oauth-utils.mock.ts | 2 + 13 files changed, 329 insertions(+), 5 deletions(-) create mode 100644 apps/sim/lib/knowledge/connectors/credential-recovery.test.ts create mode 100644 apps/sim/lib/knowledge/connectors/credential-recovery.ts diff --git a/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connector-recovery.tsx b/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connector-recovery.tsx index f95d17a948b..a5f97ed0d7d 100644 --- a/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connector-recovery.tsx +++ b/apps/sim/app/workspace/[workspaceId]/knowledge/[id]/components/connectors-section/connector-recovery.tsx @@ -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' @@ -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' @@ -97,7 +103,7 @@ export function ConnectorRecovery({ } /> )} - {connector.status === 'disabled' || credentialRemoved ? ( + {connector.status === 'disabled' || reconnectRequired ? ( { 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', diff --git a/apps/sim/lib/credentials/draft-hooks.ts b/apps/sim/lib/credentials/draft-hooks.ts index 6cf0b4c9f1f..63db1dff438 100644 --- a/apps/sim/lib/credentials/draft-hooks.ts +++ b/apps/sim/lib/credentials/draft-hooks.ts @@ -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' @@ -75,6 +76,7 @@ export async function handleCreateCredentialFromDraft(params: { .where(eq(schema.credential.id, existingCredential.id)) await clearOAuthRefreshDeadFlag(accountId) + await resumeConnectorsAfterCredentialReconnect(existingCredential.id, now) recordAudit({ workspaceId: draft.workspaceId, @@ -209,6 +211,7 @@ export async function handleReconnectCredential(params: { ) await clearOAuthRefreshDeadFlag(newAccountId) + await resumeConnectorsAfterCredentialReconnect(draft.credentialId, now) recordAudit({ workspaceId, diff --git a/apps/sim/lib/credentials/organization-draft.ts b/apps/sim/lib/credentials/organization-draft.ts index 7bc0c03fd69..fac1e793a70 100644 --- a/apps/sim/lib/credentials/organization-draft.ts +++ b/apps/sim/lib/credentials/organization-draft.ts @@ -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. */ @@ -124,6 +125,7 @@ export async function completeOrganizationCredentialDraft(input: { } }) await clearOAuthRefreshDeadFlag(input.accountId) + if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(result.credentialId, now) recordAudit({ actorId: input.userId, action: result.reconnected diff --git a/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts b/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts new file mode 100644 index 00000000000..7b976bdfe8c --- /dev/null +++ b/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts @@ -0,0 +1,31 @@ +/** + * @vitest-environment node + */ +import { dbChainMockFns, resetDbChainMock } 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' + +describe('resumeConnectorsAfterCredentialReconnect', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + }) + + it('puts only the connectors the rejected credential unscheduled back on schedule', async () => { + const now = new Date('2026-09-22T20:00:00.000Z') + await resumeConnectorsAfterCredentialReconnect('credential-1', now) + expect(dbChainMockFns.set).toHaveBeenCalledWith({ + status: 'active', + lastSyncError: null, + consecutiveFailures: 0, + nextSyncAt: now, + updatedAt: now, + }) + const guard = JSON.stringify(dbChainMockFns.where.mock.calls.at(-1)) + expect(guard).toContain('knowledgeConnector.credentialId') + expect(guard).toContain('credential-1') + expect(guard).toContain('knowledgeConnector.status') + expect(guard).toContain(CREDENTIAL_REVOKED_SYNC_ERROR) + }) +}) diff --git a/apps/sim/lib/knowledge/connectors/credential-recovery.ts b/apps/sim/lib/knowledge/connectors/credential-recovery.ts new file mode 100644 index 00000000000..9e828bc5abe --- /dev/null +++ b/apps/sim/lib/knowledge/connectors/credential-recovery.ts @@ -0,0 +1,36 @@ +import { db } from '@sim/db' +import { knowledgeConnector } from '@sim/db/schema' +import { and, eq } from 'drizzle-orm' +import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits' + +/** + * Puts the connectors a reconnected credential 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 same credential is that moment: every connector still + * carrying that error is due now, with its failure count cleared. Connectors that were paused + * or disabled for another reason keep their state, and a connector that already moved on is + * left alone. + */ +export async function resumeConnectorsAfterCredentialReconnect( + credentialId: string, + now: Date +): Promise { + await db + .update(knowledgeConnector) + .set({ + status: 'active', + lastSyncError: null, + consecutiveFailures: 0, + nextSyncAt: now, + updatedAt: now, + }) + .where( + and( + eq(knowledgeConnector.credentialId, credentialId), + eq(knowledgeConnector.status, 'error'), + eq(knowledgeConnector.lastSyncError, CREDENTIAL_REVOKED_SYNC_ERROR) + ) + ) +} diff --git a/apps/sim/lib/knowledge/connectors/sync-engine.test.ts b/apps/sim/lib/knowledge/connectors/sync-engine.test.ts index a0f8f97d284..d8b3a3c2e68 100644 --- a/apps/sim/lib/knowledge/connectors/sync-engine.test.ts +++ b/apps/sim/lib/knowledge/connectors/sync-engine.test.ts @@ -3,6 +3,7 @@ */ import { authOAuthUtilsMock, + authOAuthUtilsMockFns, dbChainMockFns, drizzleOrmMock, flattenMockConditions, @@ -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, @@ -3007,6 +3009,84 @@ 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() + authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.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('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 }, diff --git a/apps/sim/lib/knowledge/connectors/sync-engine.ts b/apps/sim/lib/knowledge/connectors/sync-engine.ts index d5f5b0dbde1..51e16443a6b 100644 --- a/apps/sim/lib/knowledge/connectors/sync-engine.ts +++ b/apps/sim/lib/knowledge/connectors/sync-engine.ts @@ -57,6 +57,7 @@ import { CONNECTOR_FAILURE_BACKOFF_CAP_MINUTES, CONNECTOR_SYNC_MAX_DURATION_SECONDS, CREDENTIAL_REMOVED_SYNC_ERROR, + CREDENTIAL_REVOKED_SYNC_ERROR, connectorFailureBackoffMinutes, MAX_CONSECUTIVE_FAILURES, } from '@/lib/knowledge/connectors/sync-limits' @@ -88,6 +89,7 @@ import { import { hardDeleteDocuments } from '@/lib/knowledge/documents/service' import { getRetryAfterMs, isRateLimitError } from '@/lib/knowledge/documents/utils' import { ensureSourceVectorIndex } from '@/lib/knowledge/search/source-vector-indexes' +import { getCredentialTerminalRefreshError } from '@/lib/oauth/credential-service' import { connectorHasAuthSource } from '@/connectors/auth' import { CONNECTOR_REGISTRY } from '@/connectors/registry.server' import type { @@ -743,9 +745,26 @@ export function buildSyncSuccessUpdate( } } +/** + * A credential the source rejected outright: the refresh path recorded a terminal error for + * it, so no retry can produce a token until the credential is reauthorized. + */ +export class ConnectorCredentialRevokedError extends Error { + constructor( + readonly credentialId: string, + readonly errorCode: string + ) { + super(`Credential ${credentialId} was rejected by the source (${errorCode})`) + this.name = 'ConnectorCredentialRevokedError' + } +} + /** * Resolves the token a connector syncs with, failing loudly where the shared * resolver reports "no token" — a sync has no reconnect prompt to fall back to. + * A credential the source has rejected outright fails as + * {@link ConnectorCredentialRevokedError}, so the run can unschedule the + * connector instead of walking the failure ladder toward a retry that cannot help. */ async function resolveAccessToken( connector: { credentialId: string | null; encryptedApiKey: string | null }, @@ -770,6 +789,13 @@ async function resolveAccessToken( userId, authMode: connectorConfig.auth.mode, }) + const terminalError = + connectorConfig.auth.mode === 'oauth' && connector.credentialId + ? await getCredentialTerminalRefreshError(connector.credentialId) + : null + if (terminalError && connector.credentialId) { + throw new ConnectorCredentialRevokedError(connector.credentialId, terminalError) + } throw new Error(`Failed to obtain access token for credential ${connector.credentialId}`) } @@ -1427,6 +1453,46 @@ export async function executeSync( } } + if (error instanceof ConnectorCredentialRevokedError) { + /** + * Retrying cannot help until the credential is reauthorized, so the + * connector leaves its schedule with a reconnect prompt instead of + * climbing the failure ladder toward the same rejection. Reauthorizing + * the credential puts it back on schedule. The run itself is a skip: + * nothing about the source failed, and a sync that cannot start is not + * an incident to page on. + */ + logger.warn('Sync unscheduled: the source rejected the connector credential', { + connectorId, + credentialId: error.credentialId, + errorCode: error.errorCode, + }) + try { + await completeSyncLog(syncLogId, 'failed', result, { + errorMessage: CREDENTIAL_REVOKED_SYNC_ERROR, + }) + const landed = await writeTerminalConnectorState( + connectorId, + syncLogId, + buildSyncUnscheduledUpdate(new Date(), CREDENTIAL_REVOKED_SYNC_ERROR) + ) + if (!landed) { + logger.warn( + 'Unschedule discarded — connector was reclaimed while this run was executing', + { connectorId, syncLogId } + ) + } + } catch (recoveryError) { + logger.error('Failed to unschedule the connector', { + connectorId, + error: + getConnectorFailureDiagnostic(recoveryError)?.message ?? + toError(recoveryError).message, + }) + } + return { ...result, skipReason: 'credential_revoked' } + } + const diagnostic = getConnectorFailureDiagnostic(error) const errorMessage = diagnostic?.message ?? toError(error).message const retryAfterMs = getRetryAfterMs(error) diff --git a/apps/sim/lib/knowledge/connectors/sync-limits.ts b/apps/sim/lib/knowledge/connectors/sync-limits.ts index 5a2af1b0ab2..5843a917afd 100644 --- a/apps/sim/lib/knowledge/connectors/sync-limits.ts +++ b/apps/sim/lib/knowledge/connectors/sync-limits.ts @@ -25,6 +25,13 @@ export const MAX_CONSECUTIVE_FAILURES = 10 export const CREDENTIAL_REMOVED_SYNC_ERROR = 'Credential removed. Reconnect the connector to resume syncing.' +/** + * The error a connector carries once the source rejects its credential outright (a revoked or + * expired grant, not a passing failure); cleared by reauthorizing that credential. + */ +export const CREDENTIAL_REVOKED_SYNC_ERROR = + 'The source no longer accepts this credential. Reconnect it to resume syncing.' + /** * The error a connector carries once {@link MAX_CONSECUTIVE_FAILURES} disables it. * diff --git a/apps/sim/lib/oauth/credential-service.test.ts b/apps/sim/lib/oauth/credential-service.test.ts index 6debccfdbee..7808e1ef8cb 100644 --- a/apps/sim/lib/oauth/credential-service.test.ts +++ b/apps/sim/lib/oauth/credential-service.test.ts @@ -77,6 +77,7 @@ vi.mock('@/lib/oauth/terminal-errors', () => ({ })) import { + getCredentialTerminalRefreshError, getOAuthToken, getServiceAccountToken, refreshTokenIfNeeded, @@ -86,6 +87,7 @@ import { } from '@/lib/oauth/credential-service' import { isInstagramProvider, shouldProactivelyRefreshInstagramToken } from '@/lib/oauth/instagram' import { isMicrosoftProvider } from '@/lib/oauth/microsoft' +import { getOAuthRefreshCoordinationIdentity } from '@/lib/oauth/refresh-coordination' import { fanOutSlackTokenChain } from '@/lib/oauth/slack' import { isTerminalRefreshError, markCredentialDead } from '@/lib/oauth/terminal-errors' import { GOOGLE_SERVICE_ACCOUNT_PROVIDER_ID } from '@/lib/oauth/types' @@ -915,3 +917,44 @@ describe('Google service-account token minting', () => { expect(JSON.stringify(mocks.logger.error.mock.calls)).not.toContain('xxxx') }) }) + +describe('getCredentialTerminalRefreshError', () => { + beforeEach(() => { + vi.clearAllMocks() + resetDbChainMock() + mocks.getRecentTerminalError.mockResolvedValue(null) + }) + + it('reads the flag on the account the credential resolves to', async () => { + queueTableRows(credential, [ + { id: RAW_CREDENTIAL_ID, type: 'oauth', accountId: RAW_ACCOUNT_ID }, + ]) + queueTableRows(account, [{ providerId: 'confluence', providerAccountId: 'provider-subject' }]) + mocks.getRecentTerminalError.mockResolvedValueOnce('invalid_grant') + await expect(getCredentialTerminalRefreshError(RAW_CREDENTIAL_ID)).resolves.toBe( + 'invalid_grant' + ) + expect(mocks.getRecentTerminalError).toHaveBeenCalledWith( + getOAuthRefreshCoordinationIdentity(RAW_ACCOUNT_ID) + ) + }) + + it('reads a Slack credential on its installation, the scope its refresh is flagged under', async () => { + queueTableRows(credential, [ + { id: RAW_CREDENTIAL_ID, type: 'oauth', accountId: RAW_ACCOUNT_ID }, + ]) + queueTableRows(account, [{ providerId: 'slack', providerAccountId: 'TEXAMPLE-usr_U1' }]) + await getCredentialTerminalRefreshError(RAW_CREDENTIAL_ID) + expect(mocks.getRecentTerminalError).toHaveBeenCalledWith( + getOAuthRefreshCoordinationIdentity('slack:TEXAMPLE') + ) + }) + + it('reports nothing for a service account, which never refreshes a chain', async () => { + queueTableRows(credential, [ + { id: RAW_CREDENTIAL_ID, type: 'service_account', accountId: null }, + ]) + await expect(getCredentialTerminalRefreshError(RAW_CREDENTIAL_ID)).resolves.toBeNull() + expect(mocks.getRecentTerminalError).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/lib/oauth/credential-service.ts b/apps/sim/lib/oauth/credential-service.ts index 1ce2d736bfc..b3d882052cf 100644 --- a/apps/sim/lib/oauth/credential-service.ts +++ b/apps/sim/lib/oauth/credential-service.ts @@ -886,6 +886,43 @@ const REFRESH_LOCK_HEADROOM_MS = 15_000 const REFRESH_LOCK_TTL_SEC = Math.ceil((TOKEN_REFRESH_TIMEOUT_MS + REFRESH_LOCK_HEADROOM_MS) / 1000) const REFRESH_FOLLOWER_MAX_WAIT_MS = REFRESH_LOCK_TTL_SEC * 1000 +/** + * The raw scope one refresh coordinates on: the account row, or the installation for a Slack + * bot token, whose sibling rows all hold one chain. + */ +function refreshCoordinationScope( + accountId: string, + providerId: string, + providerAccountId: string | null | undefined +): string { + const slackTeamId = isSlackProvider(providerId) ? extractSlackTeamId(providerAccountId) : null + return slackTeamId ? `slack:${slackTeamId}` : accountId +} + +/** + * The terminal error the refresh path last recorded for a credential's account, if any. A + * refresh that the provider rejected outright (a revoked or expired grant) flags the account + * for an hour so nothing retries it; a caller that finds no token can read the flag to tell + * that outcome, which only reauthorizing resolves, from a passing failure worth retrying. + */ +export async function getCredentialTerminalRefreshError( + credentialId: string +): Promise { + const resolved = await resolveOAuthAccountId(credentialId) + if (!resolved || resolved.credentialType === 'service_account' || !resolved.accountId) return null + const [row] = await db + .select({ providerId: account.providerId, providerAccountId: account.accountId }) + .from(account) + .where(eq(account.id, resolved.accountId)) + .limit(1) + if (!row) return null + return getRecentTerminalError( + getOAuthRefreshCoordinationIdentity( + refreshCoordinationScope(resolved.accountId, row.providerId, row.providerAccountId) + ) + ) +} + interface StoredChain { accessToken: string | null accessTokenExpiresAt: Date | null @@ -933,8 +970,9 @@ async function performCoalescedRefresh({ * dead-flagged, and written per installation rather than per row. */ const slackTeamId = isSlackProvider(providerId) ? extractSlackTeamId(providerAccountId) : null - const rawScopeKey = slackTeamId ? `slack:${slackTeamId}` : accountId - const scopeKey = getOAuthRefreshCoordinationIdentity(rawScopeKey) + const scopeKey = getOAuthRefreshCoordinationIdentity( + refreshCoordinationScope(accountId, providerId, providerAccountId) + ) const logContext = { ...(requestId ? { requestId } : {}), diff --git a/packages/testing/src/mocks/auth-oauth-utils.mock.ts b/packages/testing/src/mocks/auth-oauth-utils.mock.ts index a11162377c6..d989465695e 100644 --- a/packages/testing/src/mocks/auth-oauth-utils.mock.ts +++ b/packages/testing/src/mocks/auth-oauth-utils.mock.ts @@ -35,6 +35,7 @@ export const authOAuthUtilsMockFns = { mockGetOAuthToken: vi.fn(), mockRefreshAccessTokenIfNeeded: vi.fn(), mockRefreshTokenIfNeeded: vi.fn(), + mockGetCredentialTerminalRefreshError: vi.fn(async () => null), } /** @@ -54,4 +55,5 @@ export const authOAuthUtilsMock = { getOAuthToken: authOAuthUtilsMockFns.mockGetOAuthToken, refreshAccessTokenIfNeeded: authOAuthUtilsMockFns.mockRefreshAccessTokenIfNeeded, refreshTokenIfNeeded: authOAuthUtilsMockFns.mockRefreshTokenIfNeeded, + getCredentialTerminalRefreshError: authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError, } From cbe3076d048b5cee77577cf825d947d2d2775e18 Mon Sep 17 00:00:00 2001 From: Waleed Latif Date: Tue, 22 Sep 2026 13:57:56 -0700 Subject: [PATCH 2/2] fix(knowledge): resume by reconnected account across a Slack installation, recheck the rejection before unscheduling, and fail a run that cannot record it --- apps/sim/lib/credentials/draft-hooks.test.ts | 1 + apps/sim/lib/credentials/draft-hooks.ts | 4 +- .../sim/lib/credentials/organization-draft.ts | 2 +- .../connectors/credential-recovery.test.ts | 47 ++++++++---- .../connectors/credential-recovery.ts | 39 +++++++--- .../knowledge/connectors/sync-engine.test.ts | 53 +++++++++++++- .../lib/knowledge/connectors/sync-engine.ts | 71 +++++++++++-------- apps/sim/lib/oauth/slack.ts | 3 +- 8 files changed, 163 insertions(+), 57 deletions(-) diff --git a/apps/sim/lib/credentials/draft-hooks.test.ts b/apps/sim/lib/credentials/draft-hooks.test.ts index 72ec213d3ff..57114ca55e9 100644 --- a/apps/sim/lib/credentials/draft-hooks.test.ts +++ b/apps/sim/lib/credentials/draft-hooks.test.ts @@ -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' }, diff --git a/apps/sim/lib/credentials/draft-hooks.ts b/apps/sim/lib/credentials/draft-hooks.ts index 63db1dff438..ed5bff2dc36 100644 --- a/apps/sim/lib/credentials/draft-hooks.ts +++ b/apps/sim/lib/credentials/draft-hooks.ts @@ -76,7 +76,7 @@ export async function handleCreateCredentialFromDraft(params: { .where(eq(schema.credential.id, existingCredential.id)) await clearOAuthRefreshDeadFlag(accountId) - await resumeConnectorsAfterCredentialReconnect(existingCredential.id, now) + await resumeConnectorsAfterCredentialReconnect(accountId, now) recordAudit({ workspaceId: draft.workspaceId, @@ -211,7 +211,7 @@ export async function handleReconnectCredential(params: { ) await clearOAuthRefreshDeadFlag(newAccountId) - await resumeConnectorsAfterCredentialReconnect(draft.credentialId, now) + await resumeConnectorsAfterCredentialReconnect(newAccountId, now) recordAudit({ workspaceId, diff --git a/apps/sim/lib/credentials/organization-draft.ts b/apps/sim/lib/credentials/organization-draft.ts index fac1e793a70..0bf7f948de9 100644 --- a/apps/sim/lib/credentials/organization-draft.ts +++ b/apps/sim/lib/credentials/organization-draft.ts @@ -125,7 +125,7 @@ export async function completeOrganizationCredentialDraft(input: { } }) await clearOAuthRefreshDeadFlag(input.accountId) - if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(result.credentialId, now) + if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(input.accountId, now) recordAudit({ actorId: input.userId, action: result.reconnected diff --git a/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts b/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts index 7b976bdfe8c..4b2fe6a2c21 100644 --- a/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts +++ b/apps/sim/lib/knowledge/connectors/credential-recovery.test.ts @@ -1,31 +1,52 @@ /** * @vitest-environment node */ -import { dbChainMockFns, resetDbChainMock } from '@sim/testing' +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('puts only the connectors the rejected credential unscheduled back on schedule', async () => { - const now = new Date('2026-09-22T20:00:00.000Z') - await resumeConnectorsAfterCredentialReconnect('credential-1', now) - expect(dbChainMockFns.set).toHaveBeenCalledWith({ - status: 'active', - lastSyncError: null, - consecutiveFailures: 0, - nextSyncAt: now, - updatedAt: now, - }) + 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-1') - expect(guard).toContain('knowledgeConnector.status') 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() }) }) diff --git a/apps/sim/lib/knowledge/connectors/credential-recovery.ts b/apps/sim/lib/knowledge/connectors/credential-recovery.ts index 9e828bc5abe..87381dde7bb 100644 --- a/apps/sim/lib/knowledge/connectors/credential-recovery.ts +++ b/apps/sim/lib/knowledge/connectors/credential-recovery.ts @@ -1,22 +1,40 @@ import { db } from '@sim/db' -import { knowledgeConnector } from '@sim/db/schema' -import { and, eq } from 'drizzle-orm' +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 credential had unscheduled back on their schedule. + * 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 same credential is that moment: every connector still - * carrying that error is due now, with its failure count cleared. Connectors that were paused - * or disabled for another reason keep their state, and a connector that already moved on is - * left alone. + * 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( - credentialId: string, + accountId: string, now: Date ): Promise { + 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({ @@ -28,7 +46,10 @@ export async function resumeConnectorsAfterCredentialReconnect( }) .where( and( - eq(knowledgeConnector.credentialId, credentialId), + inArray( + knowledgeConnector.credentialId, + db.select({ id: credential.id }).from(credential).where(repairedAccounts) + ), eq(knowledgeConnector.status, 'error'), eq(knowledgeConnector.lastSyncError, CREDENTIAL_REVOKED_SYNC_ERROR) ) diff --git a/apps/sim/lib/knowledge/connectors/sync-engine.test.ts b/apps/sim/lib/knowledge/connectors/sync-engine.test.ts index d8b3a3c2e68..fcdb6d43378 100644 --- a/apps/sim/lib/knowledge/connectors/sync-engine.test.ts +++ b/apps/sim/lib/knowledge/connectors/sync-engine.test.ts @@ -3041,9 +3041,10 @@ describe('executeSync heartbeats during the listing phase', () => { it('unschedules a connector whose credential the source rejected instead of retrying it', async () => { const restore = primeOAuthRunUpToToken() - authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValueOnce( - 'invalid_grant' - ) + /** 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, @@ -3070,6 +3071,52 @@ describe('executeSync heartbeats during the listing phase', () => { } }) + 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) diff --git a/apps/sim/lib/knowledge/connectors/sync-engine.ts b/apps/sim/lib/knowledge/connectors/sync-engine.ts index 51e16443a6b..d8a429696a3 100644 --- a/apps/sim/lib/knowledge/connectors/sync-engine.ts +++ b/apps/sim/lib/knowledge/connectors/sync-engine.ts @@ -1458,39 +1458,54 @@ export async function executeSync( * Retrying cannot help until the credential is reauthorized, so the * connector leaves its schedule with a reconnect prompt instead of * climbing the failure ladder toward the same rejection. Reauthorizing - * the credential puts it back on schedule. The run itself is a skip: - * nothing about the source failed, and a sync that cannot start is not - * an incident to page on. + * the credential puts it back on schedule, and a reauthorization that + * landed while this run was failing has already cleared the rejection: + * that run takes the ordinary ladder below, so its next attempt uses the + * repaired chain rather than leaving a repaired connector unscheduled. + * The unscheduled run itself is a skip: nothing about the source failed, + * and a sync that cannot start is not an incident to page on. A run that + * cannot record the unschedule is a failure, so the runner reports it + * instead of leaving the connector locked behind a benign outcome. */ - logger.warn('Sync unscheduled: the source rejected the connector credential', { - connectorId, - credentialId: error.credentialId, - errorCode: error.errorCode, - }) - try { - await completeSyncLog(syncLogId, 'failed', result, { - errorMessage: CREDENTIAL_REVOKED_SYNC_ERROR, - }) - const landed = await writeTerminalConnectorState( + const stillRejected = await getCredentialTerminalRefreshError(error.credentialId) + if (stillRejected) { + logger.warn('Sync unscheduled: the source rejected the connector credential', { connectorId, - syncLogId, - buildSyncUnscheduledUpdate(new Date(), CREDENTIAL_REVOKED_SYNC_ERROR) - ) - if (!landed) { - logger.warn( - 'Unschedule discarded — connector was reclaimed while this run was executing', - { connectorId, syncLogId } + credentialId: error.credentialId, + errorCode: error.errorCode, + }) + try { + await completeSyncLog(syncLogId, 'failed', result, { + errorMessage: CREDENTIAL_REVOKED_SYNC_ERROR, + }) + const landed = await writeTerminalConnectorState( + connectorId, + syncLogId, + buildSyncUnscheduledUpdate(new Date(), CREDENTIAL_REVOKED_SYNC_ERROR) ) - } - } catch (recoveryError) { - logger.error('Failed to unschedule the connector', { - connectorId, - error: + if (!landed) { + logger.warn( + 'Unschedule discarded — connector was reclaimed while this run was executing', + { connectorId, syncLogId } + ) + } + return { ...result, skipReason: 'credential_revoked' } + } catch (recoveryError) { + const recoveryMessage = getConnectorFailureDiagnostic(recoveryError)?.message ?? - toError(recoveryError).message, - }) + toError(recoveryError).message + logger.error('Failed to unschedule the connector', { + connectorId, + error: recoveryMessage, + }) + result.error = recoveryMessage + return result + } } - return { ...result, skipReason: 'credential_revoked' } + logger.info('Credential reauthorized during the run; the retry uses the repaired chain', { + connectorId, + credentialId: error.credentialId, + }) } const diagnostic = getConnectorFailureDiagnostic(error) diff --git a/apps/sim/lib/oauth/slack.ts b/apps/sim/lib/oauth/slack.ts index 9c69ed3c887..e5f149d99c3 100644 --- a/apps/sim/lib/oauth/slack.ts +++ b/apps/sim/lib/oauth/slack.ts @@ -30,7 +30,8 @@ export function extractSlackTeamId(externalAccountId: string | null | undefined) return match ? match[1] : null } -function installationFilter(teamId: string) { +/** The account rows of one Slack installation: every bot token row of the team. */ +export function installationFilter(teamId: string) { return and(eq(account.providerId, 'slack'), like(account.accountId, `${teamId}-%`)) }