Skip to content

Commit e632415

Browse files
fix(knowledge): make connector sync recovery durable
1 parent 335c015 commit e632415

4 files changed

Lines changed: 49 additions & 8 deletions

File tree

‎apps/sim/app/api/knowledge/connectors/sync/route.ts‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -102,7 +102,11 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
102102
throw new Error(`Connector ${connector.id} is missing workspace billing context`)
103103
}
104104
const billingAttribution = await resolveSystemBillingAttribution(connector.workspaceId)
105-
await dispatchSync(connector.id, { billingAttribution, requestId })
105+
await dispatchSync(connector.id, {
106+
billingAttribution,
107+
requestId,
108+
requireRunnable: true,
109+
})
106110
} catch (error) {
107111
logger.error(`[${requestId}] Failed to dispatch sync for connector ${connector.id}`, error)
108112
}

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

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -882,7 +882,12 @@ export async function executeSync(
882882
.returning({ id: knowledgeConnector.id })
883883

884884
if (lockResult.length === 0) {
885-
logger.info('Sync already in progress, skipping', { connectorId })
885+
logger.info(
886+
options.requireRunnable
887+
? 'Connector is not runnable or sync is already in progress, skipping'
888+
: 'Sync already in progress, skipping',
889+
{ connectorId }
890+
)
886891
return result
887892
}
888893

‎apps/sim/lib/knowledge/orchestration/connectors.test.ts‎

Lines changed: 31 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -254,10 +254,20 @@ describe('performUpdateKnowledgeConnector', () => {
254254

255255
it('queues synchronization after replacing an active connector source', async () => {
256256
dbChainMockFns.limit.mockResolvedValueOnce([
257-
{ id: 'conn-1', connectorType: 'notion', status: 'active' },
257+
{
258+
id: 'conn-1',
259+
connectorType: 'notion',
260+
status: 'active',
261+
syncIntervalMinutes: 0,
262+
},
258263
])
259264
dbChainMockFns.returning.mockResolvedValueOnce([
260-
{ id: 'conn-1', connectorType: 'notion', status: 'active' },
265+
{
266+
id: 'conn-1',
267+
connectorType: 'notion',
268+
status: 'active',
269+
syncIntervalMinutes: 0,
270+
},
261271
])
262272

263273
const outcome = await performUpdateKnowledgeConnector({
@@ -278,12 +288,22 @@ describe('performUpdateKnowledgeConnector', () => {
278288
})
279289
})
280290

281-
it('reports a queue failure after replacing an active connector source', async () => {
291+
it('reports a queue failure and leaves the source sync due for retry', async () => {
282292
dbChainMockFns.limit.mockResolvedValueOnce([
283-
{ id: 'conn-1', connectorType: 'notion', status: 'active' },
293+
{
294+
id: 'conn-1',
295+
connectorType: 'notion',
296+
status: 'active',
297+
syncIntervalMinutes: 0,
298+
},
284299
])
285300
dbChainMockFns.returning.mockResolvedValueOnce([
286-
{ id: 'conn-1', connectorType: 'notion', status: 'active' },
301+
{
302+
id: 'conn-1',
303+
connectorType: 'notion',
304+
status: 'active',
305+
syncIntervalMinutes: 0,
306+
},
287307
])
288308
mockDispatchSync.mockRejectedValueOnce(new Error('queue unavailable'))
289309

@@ -302,6 +322,12 @@ describe('performUpdateKnowledgeConnector', () => {
302322
error: 'queue unavailable',
303323
})
304324
expect(dbChainMockFns.update).toHaveBeenCalledOnce()
325+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
326+
expect.objectContaining({
327+
sourceConfig: { database: 'next' },
328+
nextSyncAt: expect.any(Date),
329+
})
330+
)
305331
expect(mockDispatchSync).toHaveBeenCalledOnce()
306332
})
307333

‎apps/sim/lib/knowledge/orchestration/connectors.ts‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -483,7 +483,10 @@ export async function performUpdateKnowledgeConnector(
483483
}
484484
}
485485

486-
const values: Partial<typeof knowledgeConnector.$inferInsert> = { updatedAt: new Date() }
486+
const updateTimestamp = new Date()
487+
const values: Partial<typeof knowledgeConnector.$inferInsert> = {
488+
updatedAt: updateTimestamp,
489+
}
487490
if (updates.sourceConfig !== undefined) {
488491
values.sourceConfig = updates.sourceConfig
489492
}
@@ -506,6 +509,9 @@ export async function performUpdateKnowledgeConnector(
506509
}
507510
}
508511
}
512+
if (shouldDispatchSourceSync) {
513+
values.nextSyncAt = updateTimestamp
514+
}
509515

510516
let updated: ConnectorRow
511517
try {

0 commit comments

Comments
 (0)