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
61 changes: 56 additions & 5 deletions apps/sim/lib/knowledge/connectors/sync-engine.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3009,6 +3009,9 @@ describe('executeSync heartbeats during the listing phase', () => {
}
)

/** A revoked grant recorded against the credential's account. */
const INVALID_GRANT = { errorCode: 'invalid_grant', providerId: 'google-drive' } as const

/** A locked OAuth connector whose token resolution the test controls. */
function primeOAuthRunUpToToken() {
const oauthConnector = {
Expand Down Expand Up @@ -3043,8 +3046,8 @@ describe('executeSync heartbeats during the listing phase', () => {
const restore = primeOAuthRunUpToToken()
/** Rejected at token resolution and still rejected when the run records its outcome. */
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce(INVALID_GRANT)
.mockResolvedValueOnce(INVALID_GRANT)
try {
const result = await executeSync('c-1', {
billingAttribution: { workspaceId: 'ws-1' } as never,
Expand Down Expand Up @@ -3075,7 +3078,7 @@ describe('executeSync heartbeats during the listing phase', () => {
const restore = primeOAuthRunUpToToken()
/** Rejected at token resolution, repaired by the time the run records its outcome. */
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce(INVALID_GRANT)
.mockResolvedValueOnce(null)
try {
const result = await executeSync('c-1', {
Expand All @@ -3098,8 +3101,8 @@ describe('executeSync heartbeats during the listing phase', () => {
const restore = primeOAuthRunUpToToken()
/** Rejected at token resolution and still rejected when the run records its outcome. */
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce('invalid_grant')
.mockResolvedValueOnce(INVALID_GRANT)
.mockResolvedValueOnce(INVALID_GRANT)
/** The terminal write fails after the lock CAS consumed the first result. */
dbChainMockFns.returning.mockReset()
dbChainMockFns.returning.mockResolvedValueOnce([
Expand All @@ -3117,6 +3120,54 @@ describe('executeSync heartbeats during the listing phase', () => {
}
})

it('unschedules a Confluence connector whose refresh token the source rejected as unauthorized_client', async () => {
const restore = primeOAuthRunUpToToken()
const rejection = { errorCode: 'unauthorized_client', providerId: 'confluence' }
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
.mockResolvedValueOnce(rejection)
.mockResolvedValueOnce(rejection)
try {
const result = await executeSync('c-1', {
billingAttribution: { workspaceId: 'ws-1' } as never,
})
expect(result.skipReason).toBe('credential_revoked')
expect(dbChainMockFns.set).toHaveBeenCalledWith(
expect.objectContaining({ nextSyncAt: null, lastSyncError: CREDENTIAL_REVOKED_SYNC_ERROR })
)
} finally {
restore()
}
})

it.each([
{ errorCode: 'invalid_client', providerId: 'google-drive' },
{ errorCode: 'bad_client_secret', providerId: 'slack' },
{ errorCode: 'invalid_client', providerId: 'confluence' },
{ errorCode: 'unauthorized_client', providerId: 'microsoft' },
])(
'keeps the failure ladder when the refresh failed on an app-registration fault: %j',
async (rejection) => {
const restore = primeOAuthRunUpToToken()
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValue(rejection)
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 })
)
expect(dbChainMockFns.set).not.toHaveBeenCalledWith(
expect.objectContaining({ lastSyncError: CREDENTIAL_REVOKED_SYNC_ERROR })
)
} finally {
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValue(null)
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)
Expand Down
30 changes: 22 additions & 8 deletions apps/sim/lib/knowledge/connectors/sync-engine.ts
Original file line number Diff line number Diff line change
Expand Up @@ -90,6 +90,7 @@ 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 { isCredentialRevocationError } from '@/lib/oauth/terminal-errors'
import { connectorHasAuthSource } from '@/connectors/auth'
import { CONNECTOR_REGISTRY } from '@/connectors/registry.server'
import type {
Expand Down Expand Up @@ -746,8 +747,8 @@ 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.
* A credential the source revoked: the refresh path recorded a revocation for it, so no retry
* can produce a token until the credential is reauthorized.
*/
export class ConnectorCredentialRevokedError extends Error {
constructor(
Expand All @@ -759,10 +760,23 @@ export class ConnectorCredentialRevokedError extends Error {
}
}

/**
* The revocation code a credential's refresh was last rejected with, if its grant is gone. An
* app-registration fault is terminal for the refresh too, but it is ours to fix and fixing it
* restores the credential, so it reads as nothing here and the connector keeps its retry ladder
* rather than waiting on an owner to reconnect a credential that was never broken.
*/
async function getCredentialRevocationError(credentialId: string): Promise<string | null> {
const rejection = await getCredentialTerminalRefreshError(credentialId)
return rejection && isCredentialRevocationError(rejection.errorCode, rejection.providerId)
? rejection.errorCode
: null
}

/**
* 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
* A credential the source has revoked fails as
* {@link ConnectorCredentialRevokedError}, so the run can unschedule the
* connector instead of walking the failure ladder toward a retry that cannot help.
*/
Expand All @@ -789,12 +803,12 @@ async function resolveAccessToken(
userId,
authMode: connectorConfig.auth.mode,
})
const terminalError =
const revocationError =
connectorConfig.auth.mode === 'oauth' && connector.credentialId
? await getCredentialTerminalRefreshError(connector.credentialId)
? await getCredentialRevocationError(connector.credentialId)
: null
if (terminalError && connector.credentialId) {
throw new ConnectorCredentialRevokedError(connector.credentialId, terminalError)
if (revocationError && connector.credentialId) {
throw new ConnectorCredentialRevokedError(connector.credentialId, revocationError)
}
throw new Error(`Failed to obtain access token for credential ${connector.credentialId}`)
}
Expand Down Expand Up @@ -1467,7 +1481,7 @@ export async function executeSync(
* cannot record the unschedule is a failure, so the runner reports it
* instead of leaving the connector locked behind a benign outcome.
*/
const stillRejected = await getCredentialTerminalRefreshError(error.credentialId)
const stillRejected = await getCredentialRevocationError(error.credentialId)
if (stillRejected) {
logger.warn('Sync unscheduled: the source rejected the connector credential', {
connectorId,
Expand Down
42 changes: 42 additions & 0 deletions apps/sim/lib/oauth/__tests__/terminal-errors.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import {
import {
clearDeadFlag,
getRecentTerminalError,
isCredentialRevocationError,
isTerminalRefreshError,
markCredentialDead,
} from '@/lib/oauth/terminal-errors'
Expand Down Expand Up @@ -73,6 +74,47 @@ describe('isTerminalRefreshError', () => {
)
})

describe('isCredentialRevocationError', () => {
it.each([
'invalid_refresh_token',
'bad_refresh_token',
'invalid_grant',
'access_denied',
'token_revoked',
])('treats %s as a revoked credential', (code) => {
expect(isCredentialRevocationError(code)).toBe(true)
})

it.each(['invalid_client', 'bad_client_secret', 'invalid_client_id', 'bad_redirect_uri'])(
'treats the app-registration fault %s as terminal but not a revocation',
(code) => {
expect(isTerminalRefreshError(code, 'confluence')).toBe(true)
expect(isCredentialRevocationError(code, 'confluence')).toBe(false)
}
)

it.each(['confluence', 'jira'])(
'treats unauthorized_client as a revocation for %s',
(providerId) => {
expect(isCredentialRevocationError('unauthorized_client', providerId)).toBe(true)
}
)

it.each([undefined, 'microsoft', 'salesforce', 'constructor', '__proto__'])(
'does not treat unauthorized_client as a revocation for %s',
(providerId) => {
expect(isCredentialRevocationError('unauthorized_client', providerId)).toBe(false)
}
)

it.each(['ratelimited', 'internal_error', undefined, null, ''])(
'returns false for %s',
(code) => {
expect(isCredentialRevocationError(code as string | undefined | null)).toBe(false)
}
)
})

describe('markCredentialDead / getRecentTerminalError / clearDeadFlag', () => {
it('roundtrips a code through Redis', async () => {
const redis = createFakeRedis()
Expand Down
15 changes: 12 additions & 3 deletions apps/sim/lib/oauth/credential-service.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -932,9 +932,10 @@ describe('getCredentialTerminalRefreshError', () => {
])
queueTableRows(account, [{ providerId: 'confluence', providerAccountId: 'provider-subject' }])
mocks.getRecentTerminalError.mockResolvedValueOnce('invalid_grant')
await expect(getCredentialTerminalRefreshError(RAW_CREDENTIAL_ID)).resolves.toBe(
'invalid_grant'
)
await expect(getCredentialTerminalRefreshError(RAW_CREDENTIAL_ID)).resolves.toEqual({
errorCode: 'invalid_grant',
providerId: 'confluence',
})
expect(mocks.getRecentTerminalError).toHaveBeenCalledWith(
getOAuthRefreshCoordinationIdentity(RAW_ACCOUNT_ID)
)
Expand All @@ -951,6 +952,14 @@ describe('getCredentialTerminalRefreshError', () => {
)
})

it('reports nothing for an account with no flag', async () => {
queueTableRows(credential, [
{ id: RAW_CREDENTIAL_ID, type: 'oauth', accountId: RAW_ACCOUNT_ID },
])
queueTableRows(account, [{ providerId: 'confluence', providerAccountId: 'provider-subject' }])
await expect(getCredentialTerminalRefreshError(RAW_CREDENTIAL_ID)).resolves.toBeNull()
})

it('reports nothing for a service account, which never refreshes a chain', async () => {
queueTableRows(credential, [
{ id: RAW_CREDENTIAL_ID, type: 'service_account', accountId: null },
Expand Down
19 changes: 14 additions & 5 deletions apps/sim/lib/oauth/credential-service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -899,15 +899,23 @@ function refreshCoordinationScope(
return slackTeamId ? `slack:${slackTeamId}` : accountId
}

/** A terminal refresh rejection recorded for a credential's account. */
export interface CredentialTerminalRefreshError {
errorCode: string
/** The account's provider, which decides what some codes mean (see `isCredentialRevocationError`). */
providerId: string
}

/**
* 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.
* refresh that the provider rejected outright (a revoked grant, or a misconfigured app
* registration) flags the account for an hour so nothing retries it; a caller that finds no
* token can read the flag to tell that outcome from a passing failure worth retrying, and
* classify it with `isCredentialRevocationError` to tell whether only reauthorizing resolves it.
*/
export async function getCredentialTerminalRefreshError(
credentialId: string
): Promise<string | null> {
): Promise<CredentialTerminalRefreshError | null> {
const resolved = await resolveOAuthAccountId(credentialId)
if (!resolved || resolved.credentialType === 'service_account' || !resolved.accountId) return null
const [row] = await db
Expand All @@ -916,11 +924,12 @@ export async function getCredentialTerminalRefreshError(
.where(eq(account.id, resolved.accountId))
.limit(1)
if (!row) return null
return getRecentTerminalError(
const errorCode = await getRecentTerminalError(
getOAuthRefreshCoordinationIdentity(
refreshCoordinationScope(resolved.accountId, row.providerId, row.providerAccountId)
)
)
return errorCode ? { errorCode, providerId: row.providerId } : null
}

interface StoredChain {
Expand Down
46 changes: 37 additions & 9 deletions apps/sim/lib/oauth/terminal-errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,26 +4,37 @@ import { getRedisClient } from '@/lib/core/config/redis'

const logger = createLogger('OAuthTerminalErrors')

/** Refresh error codes that no retry can recover from: the credential stays dead until its owner reconnects. */
const TERMINAL_ERRORS = new Set<string>([
/**
* Refresh error codes that say the credential itself is gone: its grant was revoked, expired, or
* rotated out, and only its owner reconnecting restores it.
*/
const CREDENTIAL_REVOCATION_ERRORS = new Set<string>([
'invalid_refresh_token',
'bad_refresh_token',
'invalid_grant',
'access_denied',
'token_revoked',
])

/**
* Refresh error codes that say our app registration is misconfigured (a rotated client secret,
* a wrong client id or redirect URI). No retry recovers them either, but the fault is ours and a
* configuration fix restores every credential of the provider without its owner doing anything.
*/
const APP_CONFIGURATION_ERRORS = new Set<string>([
'bad_client_secret',
'invalid_client_id',
'invalid_client',
'bad_redirect_uri',
'token_revoked',
])

/**
* Codes terminal only for the providers listed. Atlassian rejects a revoked or rotated-out
* refresh token with `unauthorized_client`; elsewhere that code usually describes the app
* registration, and treating it as terminal would send every credential of the provider to
* Credential revocation codes only for the providers listed. Atlassian rejects a revoked or
* rotated-out refresh token with `unauthorized_client`; elsewhere that code usually describes the
* app registration, and treating it as terminal would send every credential of the provider to
* reauthorization over one configuration fault.
*/
const PROVIDER_TERMINAL_ERRORS: ReadonlyMap<string, ReadonlySet<string>> = new Map([
const PROVIDER_CREDENTIAL_REVOCATION_ERRORS: ReadonlyMap<string, ReadonlySet<string>> = new Map([
['confluence', new Set(['unauthorized_client'])],
['jira', new Set(['unauthorized_client'])],
])
Expand All @@ -34,13 +45,30 @@ function deadKey(accountId: string): string {
return `oauth:dead:${accountId}`
}

/**
* Whether a refresh error code means the credential's own grant is gone, so that only its owner
* reconnecting restores it. App-registration faults are terminal for a refresh but not a
* revocation: fixing the configuration restores the credential with no action from its owner.
*/
export function isCredentialRevocationError(
code: string | undefined | null,
providerId?: string
): boolean {
if (!code) return false
if (CREDENTIAL_REVOCATION_ERRORS.has(code)) return true
return (
providerId !== undefined &&
(PROVIDER_CREDENTIAL_REVOCATION_ERRORS.get(providerId)?.has(code) ?? false)
)
}

/** Whether no retry of a refresh can recover from the error code, whoever's fault it is. */
export function isTerminalRefreshError(
code: string | undefined | null,
providerId?: string
): boolean {
if (!code) return false
if (TERMINAL_ERRORS.has(code)) return true
return providerId !== undefined && (PROVIDER_TERMINAL_ERRORS.get(providerId)?.has(code) ?? false)
return APP_CONFIGURATION_ERRORS.has(code) || isCredentialRevocationError(code, providerId)
}

export async function markCredentialDead(accountId: string, code: string): Promise<void> {
Expand Down
4 changes: 3 additions & 1 deletion packages/testing/src/mocks/auth-oauth-utils.mock.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,9 @@ export const authOAuthUtilsMockFns = {
mockGetOAuthToken: vi.fn(),
mockRefreshAccessTokenIfNeeded: vi.fn(),
mockRefreshTokenIfNeeded: vi.fn(),
mockGetCredentialTerminalRefreshError: vi.fn(async () => null),
mockGetCredentialTerminalRefreshError: vi.fn(
async (): Promise<{ errorCode: string; providerId: string } | null> => null
),
}

/**
Expand Down
Loading