Skip to content

Commit cbe3076

Browse files
committed
fix(knowledge): resume by reconnected account across a Slack installation, recheck the rejection before unscheduling, and fail a run that cannot record it
1 parent 5ccf8c1 commit cbe3076

8 files changed

Lines changed: 163 additions & 57 deletions

File tree

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,7 @@ describe('handleReconnectCredential', () => {
103103
{ id: 'credential-1', accountId: null, displayName: 'Renamed Gmail' },
104104
])
105105
queueTableRows(schemaMock.credential, [])
106+
queueTableRows(schemaMock.account, [{ providerId: 'gmail', accountId: 'subject-new' }])
106107

107108
await handleReconnectCredential({
108109
draft: { credentialId: 'credential-1' },

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -76,7 +76,7 @@ export async function handleCreateCredentialFromDraft(params: {
7676
.where(eq(schema.credential.id, existingCredential.id))
7777

7878
await clearOAuthRefreshDeadFlag(accountId)
79-
await resumeConnectorsAfterCredentialReconnect(existingCredential.id, now)
79+
await resumeConnectorsAfterCredentialReconnect(accountId, now)
8080

8181
recordAudit({
8282
workspaceId: draft.workspaceId,
@@ -211,7 +211,7 @@ export async function handleReconnectCredential(params: {
211211
)
212212

213213
await clearOAuthRefreshDeadFlag(newAccountId)
214-
await resumeConnectorsAfterCredentialReconnect(draft.credentialId, now)
214+
await resumeConnectorsAfterCredentialReconnect(newAccountId, now)
215215

216216
recordAudit({
217217
workspaceId,

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -125,7 +125,7 @@ export async function completeOrganizationCredentialDraft(input: {
125125
}
126126
})
127127
await clearOAuthRefreshDeadFlag(input.accountId)
128-
if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(result.credentialId, now)
128+
if (result.reconnected) await resumeConnectorsAfterCredentialReconnect(input.accountId, now)
129129
recordAudit({
130130
actorId: input.userId,
131131
action: result.reconnected
Lines changed: 34 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,31 +1,52 @@
11
/**
22
* @vitest-environment node
33
*/
4-
import { dbChainMockFns, resetDbChainMock } from '@sim/testing'
4+
import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing'
55
import { beforeEach, describe, expect, it, vi } from 'vitest'
66
import { resumeConnectorsAfterCredentialReconnect } from '@/lib/knowledge/connectors/credential-recovery'
77
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
88

9+
const RESUMED = {
10+
status: 'active',
11+
lastSyncError: null,
12+
consecutiveFailures: 0,
13+
}
14+
915
describe('resumeConnectorsAfterCredentialReconnect', () => {
16+
const now = new Date('2026-09-22T20:00:00.000Z')
17+
1018
beforeEach(() => {
1119
vi.clearAllMocks()
1220
resetDbChainMock()
1321
})
1422

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-
})
23+
it('resumes the connectors of every credential on the reconnected account', async () => {
24+
queueTableRows(schemaMock.account, [
25+
{ providerId: 'confluence', providerAccountId: 'subject-1' },
26+
])
27+
await resumeConnectorsAfterCredentialReconnect('account-1', now)
28+
expect(dbChainMockFns.set).toHaveBeenCalledWith({ ...RESUMED, nextSyncAt: now, updatedAt: now })
2529
const guard = JSON.stringify(dbChainMockFns.where.mock.calls.at(-1))
2630
expect(guard).toContain('knowledgeConnector.credentialId')
27-
expect(guard).toContain('credential-1')
28-
expect(guard).toContain('knowledgeConnector.status')
2931
expect(guard).toContain(CREDENTIAL_REVOKED_SYNC_ERROR)
32+
const credentials = JSON.stringify(dbChainMockFns.where.mock.calls)
33+
expect(credentials).toContain('"left":"credential.accountId","right":"account-1"')
34+
})
35+
36+
it('resumes across the Slack installation, whose sibling accounts share the repaired chain', async () => {
37+
queueTableRows(schemaMock.account, [
38+
{ providerId: 'slack', providerAccountId: 'TEXAMPLE-usr_U1' },
39+
])
40+
await resumeConnectorsAfterCredentialReconnect('account-1', now)
41+
expect(dbChainMockFns.set).toHaveBeenCalledWith({ ...RESUMED, nextSyncAt: now, updatedAt: now })
42+
const conditions = JSON.stringify(dbChainMockFns.where.mock.calls)
43+
expect(conditions).toContain('"pattern":"TEXAMPLE-%"')
44+
expect(conditions).not.toContain('"left":"credential.accountId","right":"account-1"')
45+
})
46+
47+
it('does nothing for an account that no longer exists', async () => {
48+
queueTableRows(schemaMock.account, [])
49+
await resumeConnectorsAfterCredentialReconnect('account-gone', now)
50+
expect(dbChainMockFns.update).not.toHaveBeenCalled()
3051
})
3152
})

‎apps/sim/lib/knowledge/connectors/credential-recovery.ts‎

Lines changed: 30 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,22 +1,40 @@
11
import { db } from '@sim/db'
2-
import { knowledgeConnector } from '@sim/db/schema'
3-
import { and, eq } from 'drizzle-orm'
2+
import { account, credential, knowledgeConnector } from '@sim/db/schema'
3+
import { and, eq, inArray } from 'drizzle-orm'
44
import { CREDENTIAL_REVOKED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
5+
import { extractSlackTeamId, installationFilter, isSlackProvider } from '@/lib/oauth/slack'
56

67
/**
7-
* Puts the connectors a reconnected credential had unscheduled back on their schedule.
8+
* Puts the connectors a reconnected account had unscheduled back on their schedule.
89
*
910
* A sync that finds its credential rejected by the source leaves the connector unscheduled
1011
* 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.
12+
* authorizes again. Reauthorizing the account is that moment: every connector on a credential
13+
* of that account still carrying the error is due now, with its failure count cleared. A Slack
14+
* reauthorization repairs the installation's shared token chain, so the connectors on every
15+
* credential of the installation's sibling accounts are due as well. Connectors paused or
16+
* disabled for another reason keep their state, and a connector that already moved on is left
17+
* alone.
1518
*/
1619
export async function resumeConnectorsAfterCredentialReconnect(
17-
credentialId: string,
20+
accountId: string,
1821
now: Date
1922
): Promise<void> {
23+
const [reconnected] = await db
24+
.select({ providerId: account.providerId, providerAccountId: account.accountId })
25+
.from(account)
26+
.where(eq(account.id, accountId))
27+
.limit(1)
28+
if (!reconnected) return
29+
const slackTeamId = isSlackProvider(reconnected.providerId)
30+
? extractSlackTeamId(reconnected.providerAccountId)
31+
: null
32+
const repairedAccounts = slackTeamId
33+
? inArray(
34+
credential.accountId,
35+
db.select({ id: account.id }).from(account).where(installationFilter(slackTeamId))
36+
)
37+
: eq(credential.accountId, accountId)
2038
await db
2139
.update(knowledgeConnector)
2240
.set({
@@ -28,7 +46,10 @@ export async function resumeConnectorsAfterCredentialReconnect(
2846
})
2947
.where(
3048
and(
31-
eq(knowledgeConnector.credentialId, credentialId),
49+
inArray(
50+
knowledgeConnector.credentialId,
51+
db.select({ id: credential.id }).from(credential).where(repairedAccounts)
52+
),
3253
eq(knowledgeConnector.status, 'error'),
3354
eq(knowledgeConnector.lastSyncError, CREDENTIAL_REVOKED_SYNC_ERROR)
3455
)

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

Lines changed: 50 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3041,9 +3041,10 @@ describe('executeSync heartbeats during the listing phase', () => {
30413041

30423042
it('unschedules a connector whose credential the source rejected instead of retrying it', async () => {
30433043
const restore = primeOAuthRunUpToToken()
3044-
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValueOnce(
3045-
'invalid_grant'
3046-
)
3044+
/** Rejected at token resolution and still rejected when the run records its outcome. */
3045+
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
3046+
.mockResolvedValueOnce('invalid_grant')
3047+
.mockResolvedValueOnce('invalid_grant')
30473048
try {
30483049
const result = await executeSync('c-1', {
30493050
billingAttribution: { workspaceId: 'ws-1' } as never,
@@ -3070,6 +3071,52 @@ describe('executeSync heartbeats during the listing phase', () => {
30703071
}
30713072
})
30723073

3074+
it('takes the failure ladder when the credential was reauthorized while the run was failing', async () => {
3075+
const restore = primeOAuthRunUpToToken()
3076+
/** Rejected at token resolution, repaired by the time the run records its outcome. */
3077+
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
3078+
.mockResolvedValueOnce('invalid_grant')
3079+
.mockResolvedValueOnce(null)
3080+
try {
3081+
const result = await executeSync('c-1', {
3082+
billingAttribution: { workspaceId: 'ws-1' } as never,
3083+
})
3084+
expect(result.skipReason).toBeUndefined()
3085+
expect(result.error).toContain('rejected by the source')
3086+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
3087+
expect.objectContaining({ status: 'error', consecutiveFailures: 1 })
3088+
)
3089+
expect(dbChainMockFns.set).not.toHaveBeenCalledWith(
3090+
expect.objectContaining({ lastSyncError: CREDENTIAL_REVOKED_SYNC_ERROR })
3091+
)
3092+
} finally {
3093+
restore()
3094+
}
3095+
})
3096+
3097+
it('reports a run that could not record the unschedule as failed, not skipped', async () => {
3098+
const restore = primeOAuthRunUpToToken()
3099+
/** Rejected at token resolution and still rejected when the run records its outcome. */
3100+
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError
3101+
.mockResolvedValueOnce('invalid_grant')
3102+
.mockResolvedValueOnce('invalid_grant')
3103+
/** The terminal write fails after the lock CAS consumed the first result. */
3104+
dbChainMockFns.returning.mockReset()
3105+
dbChainMockFns.returning.mockResolvedValueOnce([
3106+
{ ...CONNECTOR, connectorType: 'oauth', credentialId: 'cred-1', accessMode: 'workspace' },
3107+
])
3108+
dbChainMockFns.returning.mockRejectedValueOnce(new Error('connection reset'))
3109+
try {
3110+
const result = await executeSync('c-1', {
3111+
billingAttribution: { workspaceId: 'ws-1' } as never,
3112+
})
3113+
expect(result.skipReason).toBeUndefined()
3114+
expect(result.error).toContain('connection reset')
3115+
} finally {
3116+
restore()
3117+
}
3118+
})
3119+
30733120
it('keeps the failure ladder for a credential that resolved no token without a terminal error', async () => {
30743121
const restore = primeOAuthRunUpToToken()
30753122
authOAuthUtilsMockFns.mockGetCredentialTerminalRefreshError.mockResolvedValueOnce(null)

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

Lines changed: 43 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -1458,39 +1458,54 @@ export async function executeSync(
14581458
* Retrying cannot help until the credential is reauthorized, so the
14591459
* connector leaves its schedule with a reconnect prompt instead of
14601460
* 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.
1461+
* the credential puts it back on schedule, and a reauthorization that
1462+
* landed while this run was failing has already cleared the rejection:
1463+
* that run takes the ordinary ladder below, so its next attempt uses the
1464+
* repaired chain rather than leaving a repaired connector unscheduled.
1465+
* The unscheduled run itself is a skip: nothing about the source failed,
1466+
* and a sync that cannot start is not an incident to page on. A run that
1467+
* cannot record the unschedule is a failure, so the runner reports it
1468+
* instead of leaving the connector locked behind a benign outcome.
14641469
*/
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(
1470+
const stillRejected = await getCredentialTerminalRefreshError(error.credentialId)
1471+
if (stillRejected) {
1472+
logger.warn('Sync unscheduled: the source rejected the connector credential', {
14751473
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 }
1474+
credentialId: error.credentialId,
1475+
errorCode: error.errorCode,
1476+
})
1477+
try {
1478+
await completeSyncLog(syncLogId, 'failed', result, {
1479+
errorMessage: CREDENTIAL_REVOKED_SYNC_ERROR,
1480+
})
1481+
const landed = await writeTerminalConnectorState(
1482+
connectorId,
1483+
syncLogId,
1484+
buildSyncUnscheduledUpdate(new Date(), CREDENTIAL_REVOKED_SYNC_ERROR)
14831485
)
1484-
}
1485-
} catch (recoveryError) {
1486-
logger.error('Failed to unschedule the connector', {
1487-
connectorId,
1488-
error:
1486+
if (!landed) {
1487+
logger.warn(
1488+
'Unschedule discarded — connector was reclaimed while this run was executing',
1489+
{ connectorId, syncLogId }
1490+
)
1491+
}
1492+
return { ...result, skipReason: 'credential_revoked' }
1493+
} catch (recoveryError) {
1494+
const recoveryMessage =
14891495
getConnectorFailureDiagnostic(recoveryError)?.message ??
1490-
toError(recoveryError).message,
1491-
})
1496+
toError(recoveryError).message
1497+
logger.error('Failed to unschedule the connector', {
1498+
connectorId,
1499+
error: recoveryMessage,
1500+
})
1501+
result.error = recoveryMessage
1502+
return result
1503+
}
14921504
}
1493-
return { ...result, skipReason: 'credential_revoked' }
1505+
logger.info('Credential reauthorized during the run; the retry uses the repaired chain', {
1506+
connectorId,
1507+
credentialId: error.credentialId,
1508+
})
14941509
}
14951510

14961511
const diagnostic = getConnectorFailureDiagnostic(error)

‎apps/sim/lib/oauth/slack.ts‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,8 @@ export function extractSlackTeamId(externalAccountId: string | null | undefined)
3030
return match ? match[1] : null
3131
}
3232

33-
function installationFilter(teamId: string) {
33+
/** The account rows of one Slack installation: every bot token row of the team. */
34+
export function installationFilter(teamId: string) {
3435
return and(eq(account.providerId, 'slack'), like(account.accountId, `${teamId}-%`))
3536
}
3637

0 commit comments

Comments
 (0)