Skip to content

Commit 35a1c58

Browse files
committed
fix(knowledge): count a connector's documents before the completion locks and index the live set
1 parent 0e6fd98 commit 35a1c58

6 files changed

Lines changed: 28570 additions & 14 deletions

File tree

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

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2252,6 +2252,61 @@ describe('completeSuccessfulSync', () => {
22522252
}
22532253
)
22542254

2255+
it('counts the documents before taking the completion locks', async () => {
2256+
const { completeSuccessfulSync } = await import('@/lib/knowledge/connectors/sync-engine')
2257+
queueTableRows(schemaMock.knowledgeBase, [{ id: 'kb-1' }])
2258+
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-1' }])
2259+
queueTableRows(schemaMock.document, [{ count: 4 }])
2260+
dbChainMockFns.returning
2261+
.mockResolvedValueOnce([])
2262+
.mockResolvedValueOnce([{ id: 'log-1' }])
2263+
.mockResolvedValueOnce([{ id: 'c-1' }])
2264+
2265+
await expect(completeSuccessfulSync('c-1', 'kb-1', 'log-1', 60, RESULT, null)).resolves.toBe(
2266+
true
2267+
)
2268+
2269+
const countOrder = dbChainMockFns.from.mock.invocationCallOrder[0]
2270+
const transactionOrder = dbChainMockFns.transaction.mock.invocationCallOrder[0]
2271+
expect(dbChainMockFns.from.mock.calls[0][0]).toBe(schemaMock.document)
2272+
expect(countOrder).toBeLessThan(transactionOrder)
2273+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
2274+
expect.objectContaining({ status: 'active', lastSyncDocCount: 4 })
2275+
)
2276+
})
2277+
2278+
it('keeps the previous document count when the count cannot be read', async () => {
2279+
const { completeSuccessfulSync } = await import('@/lib/knowledge/connectors/sync-engine')
2280+
queueTableRows(schemaMock.knowledgeBase, [{ id: 'kb-1' }])
2281+
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'c-1' }])
2282+
dbChainMockFns.where.mockImplementationOnce(() =>
2283+
Promise.reject(
2284+
new DrizzleQueryError(
2285+
'select private SQL',
2286+
[],
2287+
Object.assign(new Error('canceling statement due to statement timeout'), {
2288+
code: '57014',
2289+
})
2290+
)
2291+
)
2292+
)
2293+
dbChainMockFns.returning
2294+
.mockResolvedValueOnce([])
2295+
.mockResolvedValueOnce([{ id: 'log-1' }])
2296+
.mockResolvedValueOnce([{ id: 'c-1' }])
2297+
2298+
await expect(completeSuccessfulSync('c-1', 'kb-1', 'log-1', 60, RESULT, null)).resolves.toBe(
2299+
true
2300+
)
2301+
2302+
const successWrite = dbChainMockFns.set.mock.calls
2303+
.map(([value]) => value as Record<string, unknown>)
2304+
.find((value) => value.status === 'active')
2305+
expect(successWrite).toBeDefined()
2306+
expect(successWrite).not.toHaveProperty('lastSyncDocCount')
2307+
expect(successWrite).toMatchObject({ consecutiveFailures: 0 })
2308+
})
2309+
22552310
it('commits the completed log and connector state in one guarded transaction', async () => {
22562311
const { completeSuccessfulSync } = await import('@/lib/knowledge/connectors/sync-engine')
22572312

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

Lines changed: 30 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -393,6 +393,18 @@ export async function completeSuccessfulSync(
393393
: null
394394
const completionNotice =
395395
[directoryNotice, listingNotice, contentNotice].filter(Boolean).join('\n') || null
396+
/**
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.
400+
*/
401+
const actualDocCount = await countLiveConnectorDocuments(connectorId).catch((error: unknown) => {
402+
logger.warn('Could not count connector documents; keeping the previous count', {
403+
connectorId,
404+
diagnostic: getConnectorFailureDiagnostic(error),
405+
})
406+
return null
407+
})
396408
try {
397409
const completed = await db.transaction(async (tx) => {
398410
const [lockedKnowledgeBase] = await tx
@@ -424,18 +436,6 @@ export async function completeSuccessfulSync(
424436
})
425437
}
426438

427-
const [{ count: actualDocCount }] = await tx
428-
.select({ count: sql<number>`count(*)::int` })
429-
.from(document)
430-
.where(
431-
and(
432-
eq(document.connectorId, connectorId),
433-
eq(document.userExcluded, false),
434-
isNull(document.archivedAt),
435-
isNull(document.deletedAt)
436-
)
437-
)
438-
439439
const now = new Date()
440440
const [closedLog] = await tx
441441
.update(knowledgeConnectorSyncLog)
@@ -703,9 +703,25 @@ export function buildSyncCapacityUpdate(
703703
* `consecutiveFailures` still resets: a held pass is a healthy sync that declined
704704
* to delete, not a failure, and marking it broken would stop it syncing at all.
705705
*/
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+
}
721+
706722
export function buildSyncSuccessUpdate(
707723
now: Date,
708-
actualDocCount: number,
724+
actualDocCount: number | null,
709725
nextSyncAt: Date | null,
710726
holdNotice: string | null,
711727
advanceLastSyncAt = true
@@ -714,7 +730,7 @@ export function buildSyncSuccessUpdate(
714730
status: 'active' as const,
715731
...(advanceLastSyncAt ? { lastSyncAt: now } : {}),
716732
lastSyncError: holdNotice,
717-
lastSyncDocCount: actualDocCount,
733+
...(actualDocCount === null ? {} : { lastSyncDocCount: actualDocCount }),
718734
nextSyncAt,
719735
consecutiveFailures: 0,
720736
// Releases the lock so a stale token can never match a later run, and closes
Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
-- Every sync completion counts the live documents a connector owns. The reconciliation index keeps
2+
-- tombstones on purpose, so without a live-only index that count walks the connector's heap tuples
3+
-- and a large connector cannot finish inside the statement timeout.
4+
COMMIT;--> statement-breakpoint
5+
SET lock_timeout = 0;--> statement-breakpoint
6+
-- migration-safe: replay replaces only this new index to recover an interrupted concurrent build; existing indexes remain available.
7+
DROP INDEX CONCURRENTLY IF EXISTS "doc_connector_live_idx";--> statement-breakpoint
8+
CREATE INDEX CONCURRENTLY IF NOT EXISTS "doc_connector_live_idx" ON "document" USING btree ("connector_id") WHERE "user_excluded" = false AND "archived_at" IS NULL AND "deleted_at" IS NULL;--> statement-breakpoint
9+
SET lock_timeout = '5s';

0 commit comments

Comments
 (0)