Skip to content

Commit 81d1ae9

Browse files
committed
fix(knowledge): share the unscheduled connector write, scope credential removal to content-engine modes, and keep both sync holds visible
1 parent 6a0469e commit 81d1ae9

8 files changed

Lines changed: 121 additions & 96 deletions

File tree

‎apps/sim/background/knowledge-connector-sync.ts‎

Lines changed: 7 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -27,17 +27,6 @@ export function classifyConnectorSyncResult(result: SyncResult): ConnectorSyncTa
2727
return 'completed'
2828
}
2929

30-
function formatConnectorSyncFailure(
31-
connectorId: string,
32-
result: SyncResult,
33-
outcome: Extract<ConnectorSyncTaskOutcome, 'partial' | 'failed'>
34-
): string {
35-
if (outcome === 'failed') {
36-
return `Connector sync failed for ${connectorId}: ${result.error}`
37-
}
38-
return `Connector sync partially failed for ${connectorId}: ${result.docsFailed} source failures, ${result.processingDispatch.failed} dispatch failures`
39-
}
40-
4130
export async function executeConnectorSyncJob(payload: unknown) {
4231
const {
4332
connectorId,
@@ -60,9 +49,10 @@ export async function executeConnectorSyncJob(payload: unknown) {
6049
dispatchToken,
6150
})
6251

52+
const outcome = classifyConnectorSyncResult(result)
6353
logger.info(`[${requestId}] Connector sync completed`, {
6454
connectorId,
65-
outcome: classifyConnectorSyncResult(result),
55+
outcome,
6656
deferred: result.deferred,
6757
added: result.docsAdded,
6858
updated: result.docsUpdated,
@@ -75,7 +65,6 @@ export async function executeConnectorSyncJob(payload: unknown) {
7565
processingDispatchFailed: result.processingDispatch.failed,
7666
})
7767

78-
const outcome = classifyConnectorSyncResult(result)
7968
if (outcome === 'failed') {
8069
/**
8170
* `executeSync` has already persisted its terminal state, and retrying this
@@ -85,10 +74,12 @@ export async function executeConnectorSyncJob(payload: unknown) {
8574
* connector pass replays them, its dispatch failures stay eligible for the
8675
* stuck-document sweep, and the outcome rides on the return value.
8776
*/
88-
throw new AbortTaskRunError(formatConnectorSyncFailure(connectorId, result, 'failed'))
77+
throw new AbortTaskRunError(`Connector sync failed for ${connectorId}: ${result.error}`)
8978
}
90-
if (outcome === 'partial') {
91-
logger.warn(`[${requestId}] ${formatConnectorSyncFailure(connectorId, result, 'partial')}`)
79+
if (outcome === 'partial' && (result.docsFailed > 0 || result.processingDispatch.failed > 0)) {
80+
logger.warn(
81+
`[${requestId}] Connector sync partially failed for ${connectorId}: ${result.docsFailed} source failures, ${result.processingDispatch.failed} dispatch failures`
82+
)
9283
}
9384

9485
return {

‎apps/sim/connectors/auth.ts‎

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,11 @@ export function connectorHasAuthSource(
3434
auth: ConnectorAuthConfig,
3535
connector: { credentialId: string | null; encryptedApiKey: string | null }
3636
): boolean {
37-
if (auth.mode === 'apiKey') return true
38-
if (connector.credentialId) return true
39-
return Boolean(getConnectorApiKeyConfig(auth) && connector.encryptedApiKey)
37+
return (
38+
auth.mode === 'apiKey' ||
39+
Boolean(connector.credentialId) ||
40+
Boolean(getConnectorApiKeyConfig(auth) && connector.encryptedApiKey)
41+
)
4042
}
4143

4244
/** Workspace token input supported by a connector, independent of its member OAuth method. */

‎apps/sim/lib/credentials/deletion.test.ts‎

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -204,22 +204,29 @@ describe('clearCredentialRefs', () => {
204204
const updates = capturedQueries.filter((query) =>
205205
query.sql.startsWith('update "knowledge_connector"')
206206
)
207-
expect(updates).toHaveLength(1)
207+
expect(updates).toHaveLength(2)
208208
const statement = normalizeSql(updates[0].sql)
209209
expect(statement).toContain('"credential_id" = $')
210+
expect(statement).toContain('"last_sync_error" = $')
210211
expect(statement).toContain('"next_sync_at" = $')
211212
expect(statement).toContain('"sync_lock_token" = $')
212213
expect(statement).toContain('"sync_lock_lease_at" = $')
213214
expect(statement).toContain(
214215
`CASE WHEN "knowledge_connector"."status" IN ('paused', 'disabled')`
215216
)
217+
expect(statement).toContain('"access_mode" in ($')
216218
expect(updates[0].params).toEqual(
217219
expect.arrayContaining([
218-
null,
219220
'Credential removed. Reconnect the connector to resume syncing.',
220221
'credential-target',
222+
'workspace',
223+
'admin',
221224
])
222225
)
226+
const rest = normalizeSql(updates[1].sql)
227+
expect(rest).toContain('"credential_id" = $')
228+
expect(rest).not.toContain('"last_sync_error"')
229+
expect(updates[1].params).toEqual(expect.arrayContaining([null, 'credential-target']))
223230
})
224231

225232
it('propagates database failures', async () => {

‎apps/sim/lib/credentials/deletion.ts‎

Lines changed: 18 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@ import { AuditAction, AuditResourceType, recordAudit } from '@sim/audit'
22
import { db } from '@sim/db'
33
import * as schema from '@sim/db/schema'
44
import { createLogger } from '@sim/logger'
5-
import { and, eq, notExists, or, sql } from 'drizzle-orm'
5+
import { and, eq, inArray, notExists, or, sql } from 'drizzle-orm'
66
import type { AnyPgColumn, PgTable } from 'drizzle-orm/pg-core'
77
import type { NextRequest } from 'next/server'
88
import {
@@ -11,7 +11,9 @@ import {
1111
resourceScopeFromOwner,
1212
} from '@/lib/core/resource-scope'
1313
import { resourceScopeCondition } from '@/lib/core/resource-scope.server'
14+
import { CONTENT_ENGINE_ACCESS_MODES } from '@/lib/knowledge/connectors/access-modes'
1415
import { CREDENTIAL_REMOVED_SYNC_ERROR } from '@/lib/knowledge/connectors/sync-limits'
16+
import { buildSyncUnscheduledUpdate } from '@/lib/knowledge/connectors/sync-lock'
1517
import { CREDENTIAL_SUBBLOCK_IDS } from '@/lib/workflows/persistence/utils'
1618

1719
const logger = createLogger('CredentialDeletion')
@@ -338,24 +340,29 @@ async function readWorkspaceCredentialRefs(
338340
}
339341

340342
/**
341-
* A connector whose credential is gone cannot sync until it is reconnected, so the scheduler
342-
* must stop dispatching it: `nextSyncAt` leaves the due sweep, the status and error name the
343-
* cause for the reconnect prompt, and the lock is released so a live run's terminal write cannot
344-
* resurrect a schedule (same transition as a deleted knowledge base). Paused and disabled
345-
* connectors keep their status.
343+
* A content-engine connector whose credential is gone cannot sync until it is reconnected, so
344+
* it leaves the due sweep with the reconnect error; paused and disabled connectors keep their
345+
* status. A members-mode connector only loses its optional dedicated content credential: its
346+
* member crawls keep running, so it merely drops the reference.
346347
*/
347348
async function clearInKnowledgeConnectors(credentialId: string): Promise<void> {
349+
const now = new Date()
348350
await db
349351
.update(schema.knowledgeConnector)
350352
.set({
353+
...buildSyncUnscheduledUpdate(now, CREDENTIAL_REMOVED_SYNC_ERROR),
351354
credentialId: null,
352355
status: sql`CASE WHEN ${schema.knowledgeConnector.status} IN ('paused', 'disabled') THEN ${schema.knowledgeConnector.status} ELSE 'error' END`,
353-
lastSyncError: CREDENTIAL_REMOVED_SYNC_ERROR,
354-
nextSyncAt: null,
355-
syncLockToken: null,
356-
syncLockLeaseAt: null,
357-
updatedAt: new Date(),
358356
})
357+
.where(
358+
and(
359+
eq(schema.knowledgeConnector.credentialId, credentialId),
360+
inArray(schema.knowledgeConnector.accessMode, [...CONTENT_ENGINE_ACCESS_MODES])
361+
)
362+
)
363+
await db
364+
.update(schema.knowledgeConnector)
365+
.set({ credentialId: null, updatedAt: now })
359366
.where(eq(schema.knowledgeConnector.credentialId, credentialId))
360367
}
361368

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

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -525,6 +525,33 @@ describe('content pass checkpoint intent', () => {
525525
expect(pass.checkpoint.permissionFailures).toBe(true)
526526
})
527527

528+
it('reports an incomplete listing alongside unverified permissions', async () => {
529+
sourceBody = { value: '<p>Current content</p>' }
530+
const checkpoint = {
531+
...beginListingCheckpoint({
532+
fingerprint: 'a'.repeat(64),
533+
generationId: 'prior',
534+
startedAt: new Date(0),
535+
}),
536+
listingFailures: {
537+
count: 1,
538+
samples: [
539+
{
540+
scope: 'unavailable@example.com',
541+
operation: 'gmail.threads.list',
542+
status: 400,
543+
reasons: ['failedPrecondition'],
544+
},
545+
],
546+
},
547+
}
548+
mocks.onPage.mockResolvedValue({ permissionsIncomplete: true })
549+
const { pass } = await runPass({ checkpoint, access: 'admin' })
550+
const lines = pass.holdNotice?.split('\n') ?? []
551+
expect(lines).toContain(SOURCE_PERMISSION_ERROR)
552+
expect(lines.some((line) => line.includes('unlisted documents were kept'))).toBe(true)
553+
})
554+
528555
it('clears permission failure evidence for a newly verified crawl', async () => {
529556
mocks.onPage.mockResolvedValue({ permissionsIncomplete: false })
530557
const { pass } = await runPass({ access: 'admin' })

‎apps/sim/lib/knowledge/connectors/sync-content-pass.ts‎

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -330,12 +330,18 @@ export async function runConnectorContentPass(input: ContentPassInput) {
330330
const reconciliation = checkpoint.complete
331331
? await reconcileCompletedListing(input, checkpoint, withLease)
332332
: { finished: false, notice: null }
333+
/** Unverified permissions and an incomplete listing are independent holds; an admin needs to see both, one per line. */
334+
const holdNotice =
335+
[
336+
checkpoint.permissionFailures ? SOURCE_PERMISSION_ERROR : null,
337+
reconciliation.notice ?? (checkpoint.contentFailures ? SOURCE_CONTENT_ERROR : null),
338+
]
339+
.filter(Boolean)
340+
.join('\n') || null
333341
return {
334342
checkpoint,
335343
complete: checkpoint.complete && reconciliation.finished,
336-
holdNotice: checkpoint.permissionFailures
337-
? SOURCE_PERMISSION_ERROR
338-
: (reconciliation.notice ?? (checkpoint.contentFailures ? SOURCE_CONTENT_ERROR : null)),
344+
holdNotice,
339345
hydratedCount,
340346
}
341347
}

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

Lines changed: 27 additions & 61 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ import {
6363
import {
6464
assertSyncLeaseHeldInTx,
6565
buildSyncLockAcquisition,
66+
buildSyncUnscheduledUpdate,
6667
createContentSyncLease,
6768
holdsSyncLockToken,
6869
LOCKABLE_CONNECTOR_STATUSES,
@@ -352,6 +353,22 @@ export function isContentPassIncomplete(
352353
)
353354
}
354355

356+
/** Live documents the connector owns: what the connector list shows as its document count. */
357+
async function countLiveConnectorDocuments(connectorId: string): Promise<number> {
358+
const [row] = await db
359+
.select({ count: sql<number>`count(*)::int` })
360+
.from(document)
361+
.where(
362+
and(
363+
eq(document.connectorId, connectorId),
364+
eq(document.userExcluded, false),
365+
isNull(document.archivedAt),
366+
isNull(document.deletedAt)
367+
)
368+
)
369+
return row?.count ?? 0
370+
}
371+
355372
/**
356373
* Atomically publishes the completed log and connector terminal state.
357374
*
@@ -394,9 +411,11 @@ export async function completeSuccessfulSync(
394411
const completionNotice =
395412
[directoryNotice, listingNotice, contentNotice].filter(Boolean).join('\n') || null
396413
/**
397-
* A display snapshot, taken before the completion transaction so a slow count neither holds
398-
* the connector lock nor turns a sync whose documents already landed into a failure. When it
399-
* cannot be read the previous count stands.
414+
* Read before the completion transaction so a slow count neither holds the connector lock
415+
* nor turns a sync whose documents already landed into a failure. It is also the previous
416+
* owned count the next pass's mass-deletion guard starts from, so a tally of this run's
417+
* adds and deletes would drift; when it cannot be read the previous count stands, which the
418+
* guard already treats as a floor.
400419
*/
401420
const actualDocCount = await countLiveConnectorDocuments(connectorId).catch((error: unknown) => {
402421
logger.warn('Could not count connector documents; keeping the previous count', {
@@ -569,14 +588,7 @@ async function releaseSyncLockOnDeletedConnector(
569588
): Promise<void> {
570589
await db
571590
.update(knowledgeConnector)
572-
.set({
573-
status: 'error',
574-
nextSyncAt: null,
575-
lastSyncError: 'Connector deleted during sync',
576-
syncLockToken: null,
577-
syncLockLeaseAt: null,
578-
updatedAt: new Date(),
579-
})
591+
.set(buildSyncUnscheduledUpdate(new Date(), 'Connector deleted during sync'))
580592
.where(holdsSyncLockToken(connectorId, syncLogId))
581593
}
582594

@@ -684,13 +696,8 @@ export function buildSyncCapacityUpdate(
684696
errorMessage: string
685697
) {
686698
return {
687-
status: 'error' as const,
688-
lastSyncError: errorMessage,
689-
nextSyncAt: null,
699+
...buildSyncUnscheduledUpdate(now, errorMessage),
690700
consecutiveFailures: previousFailures ?? 0,
691-
syncLockToken: null,
692-
syncLockLeaseAt: null,
693-
updatedAt: now,
694701
}
695702
}
696703

@@ -703,21 +710,6 @@ export function buildSyncCapacityUpdate(
703710
* `consecutiveFailures` still resets: a held pass is a healthy sync that declined
704711
* to delete, not a failure, and marking it broken would stop it syncing at all.
705712
*/
706-
/** Live documents the connector owns: what the connector list shows as its document count. */
707-
async function countLiveConnectorDocuments(connectorId: string): Promise<number> {
708-
const [row] = await db
709-
.select({ count: sql<number>`count(*)::int` })
710-
.from(document)
711-
.where(
712-
and(
713-
eq(document.connectorId, connectorId),
714-
eq(document.userExcluded, false),
715-
isNull(document.archivedAt),
716-
isNull(document.deletedAt)
717-
)
718-
)
719-
return row?.count ?? 0
720-
}
721713

722714
export function buildSyncSuccessUpdate(
723715
now: Date,
@@ -850,21 +842,13 @@ export async function executeSync(
850842
/**
851843
* A connector with no token source cannot succeed, and each attempt would only walk the
852844
* failure ladder and, at its end, disable a connector that merely needs reconnecting. Left
853-
* unscheduled with the reconnect error instead; the same terminal transition as a deleted
854-
* knowledge base, so a live run's own terminal write cannot revive the schedule.
845+
* unscheduled with the reconnect error instead, the same transition as a deleted knowledge base.
855846
*/
856847
if (!connectorHasAuthSource(connectorConfig.auth, connectorBeforeLock)) {
857848
logger.warn('Skipping sync: connector has no credential to authenticate with', { connectorId })
858849
await db
859850
.update(knowledgeConnector)
860-
.set({
861-
status: 'error',
862-
nextSyncAt: null,
863-
lastSyncError: CREDENTIAL_REMOVED_SYNC_ERROR,
864-
syncLockToken: null,
865-
syncLockLeaseAt: null,
866-
updatedAt: new Date(),
867-
})
851+
.set(buildSyncUnscheduledUpdate(new Date(), CREDENTIAL_REMOVED_SYNC_ERROR))
868852
.where(eq(knowledgeConnector.id, connectorId))
869853
return { ...result, skipReason: 'credential_missing' }
870854
}
@@ -890,25 +874,7 @@ export async function executeSync(
890874
)
891875
await db
892876
.update(knowledgeConnector)
893-
.set({
894-
status: 'error',
895-
nextSyncAt: null,
896-
lastSyncError: 'Knowledge base deleted',
897-
/**
898-
* Clears the lock alongside the status.
899-
*
900-
* This write runs BEFORE the lock is taken, but it is unconditional on
901-
* status, so it can land on a row a previous run left `syncing` — a run
902-
* that may still be alive. Flipping status without releasing the token
903-
* left a row that was neither locked nor reclaimable: the reaper only
904-
* looks at `syncing` rows, and the old run's terminal write could still
905-
* match its own token and resurrect a state for a knowledge base that no
906-
* longer exists. Releasing both makes the transition terminal.
907-
*/
908-
syncLockToken: null,
909-
syncLockLeaseAt: null,
910-
updatedAt: new Date(),
911-
})
877+
.set(buildSyncUnscheduledUpdate(new Date(), 'Knowledge base deleted'))
912878
.where(eq(knowledgeConnector.id, connectorId))
913879
return { ...result, skipReason: 'knowledge_base_deleted' }
914880
}

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

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,25 @@ export function holdsSyncLockToken(connectorId: string, syncLockToken: string) {
6969
/** Connector statuses a scheduler may start an automatic run from. */
7070
export const RUNNABLE_CONNECTOR_STATUSES = ['active', 'error'] as const
7171

72+
/**
73+
* A terminal, unscheduled error. The lock is released alongside the status because
74+
* this write can land on a row a previous run left `syncing` — a run that may still
75+
* be alive. Flipping status without releasing the token left a row that was neither
76+
* locked nor reclaimable: the reaper only looks at `syncing` rows, and the old run's
77+
* terminal write could still match its own token and resurrect the schedule.
78+
* Releasing both makes the transition terminal.
79+
*/
80+
export function buildSyncUnscheduledUpdate(now: Date, errorMessage: string) {
81+
return {
82+
status: 'error' as const,
83+
lastSyncError: errorMessage,
84+
nextSyncAt: null,
85+
syncLockToken: null,
86+
syncLockLeaseAt: null,
87+
updatedAt: now,
88+
}
89+
}
90+
7291
/**
7392
* The statuses a run may take the lock from.
7493
*

0 commit comments

Comments
 (0)