Skip to content

Commit cc9abe0

Browse files
committed
fix(knowledge): batch scope renewal and keep refreshed metadata
Scope renewal now gathers container pages into batches before scanning the member's stale observations, so the scan runs once per batch instead of once per source page, and an unfinished batch resumes from where it was read. Content that hydrates unchanged under a new hash also refreshes its source URL, modified time and tags, and the member sync log records how many observations renewal kept fresh.
1 parent f15fb7c commit cc9abe0

13 files changed

Lines changed: 185 additions & 18 deletions

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ export async function executeMemberSyncJob(payload: unknown) {
6868
unchanged: result.docsUnchanged,
6969
failed: result.docsFailed,
7070
observationsAdded: result.observationsAdded,
71+
observationsRenewed: result.observationsRenewed,
7172
observationsRemoved: result.observationsRemoved,
7273
tombstoned: result.docsTombstoned,
7374
resurrected: result.docsResurrected,

‎apps/sim/lib/api/contracts/knowledge/connectors.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -206,6 +206,8 @@ export const memberSyncLogDataSchema = z
206206
docsUnchanged: z.number(),
207207
docsHydratedOnce: z.number(),
208208
observationsAdded: z.number(),
209+
/** Absent from responses served by older deployments. */
210+
observationsRenewed: z.number().int().nonnegative().optional(),
209211
observationsRemoved: z.number(),
210212
docsTombstoned: z.number(),
211213
docsResurrected: z.number(),

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

Lines changed: 57 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ vi.mock('@/lib/knowledge/connectors/sync-persistence', () => ({
7171
persistSkippedDocuments: vi.fn(async () => []),
7272
persistSourceDocumentFailures: mocks.persistFailures,
7373
persistHashOnlyUpdates: vi.fn(async () => []),
74+
resolveSourceMetadataFields: vi.fn(() => ({ sourceUrl: null, sourceModifiedAt: null })),
7475
}))
7576
vi.mock('@/lib/knowledge/documents/service', () => ({
7677
isTriggerAvailable: () => true,
@@ -117,6 +118,7 @@ vi.mock('@/connectors/registry.server', () => ({
117118
},
118119
}))
119120

121+
import { ProviderCapacityDeferredError } from '@/lib/core/rate-limiter/provider-capacity-error'
120122
import {
121123
CredentialGroupCredentialCursorNotFoundError,
122124
loadScopedAccountsCredentialListContext,
@@ -126,7 +128,10 @@ import {
126128
listingFingerprint,
127129
} from '@/lib/knowledge/connectors/listing-checkpoint'
128130
import { executeMemberSync } from '@/lib/knowledge/connectors/member-sync-engine'
129-
import { SOURCE_CONTENT_ERROR } from '@/lib/knowledge/connectors/sync-limits'
131+
import {
132+
MEMBER_SCOPE_RENEWAL_PREFIX_BATCH,
133+
SOURCE_CONTENT_ERROR,
134+
} from '@/lib/knowledge/connectors/sync-limits'
130135

131136
const serviceDocument: ExternalDocument = {
132137
externalId: 'file-shared',
@@ -858,23 +863,59 @@ describe('member engine with a dedicated content credential', () => {
858863
scopeRenewalCursor: null,
859864
scopeRenewalStartedAt: null,
860865
})
866+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
867+
expect.objectContaining({ completedAt: expect.any(Date), observationsRenewed: 3 })
868+
)
861869
})
862870

863-
it('renews every page of scopes before recording the renewal', async () => {
871+
it('records renewed observations on the log of a deferred run', async () => {
872+
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
873+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'] })
874+
mocks.renew.mockResolvedValue({ renewed: 2, finished: true })
875+
mocks.list.mockRejectedValue(new ProviderCapacityDeferredError('admission_unavailable'))
876+
expect((await run()).deferred).toBeDefined()
877+
expect(dbChainMockFns.set).toHaveBeenCalledWith(
878+
expect.objectContaining({ completedAt: expect.any(Date), observationsRenewed: 2 })
879+
)
880+
})
881+
882+
it('gathers scope pages into one pass over the stale observations', async () => {
864883
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
865884
mocks.scopes
866885
.mockResolvedValueOnce({ prefixes: ['source:container-a:'], nextCursor: 'page-2' })
867886
.mockResolvedValueOnce({ prefixes: ['source:container-b:'] })
868887
mocks.renew.mockResolvedValue({ renewed: 2, finished: true })
869-
expect((await run()).observationsRenewed).toBe(4)
888+
expect((await run()).observationsRenewed).toBe(2)
870889
expect(mocks.scopes.mock.calls.map((call) => call[2])).toEqual([undefined, 'page-2'])
871890
expect(mocks.renew.mock.calls.map(([call]) => call.scopePrefixes)).toEqual([
872-
['source:container-a:'],
873-
['source:container-b:'],
891+
['source:container-a:', 'source:container-b:'],
874892
])
875893
expect(scopeRenewal()).toMatchObject({ scopeRenewedAt: expect.any(Date) })
876894
})
877895

896+
it('renews a full batch of scopes before reading more, and resumes an unfinished batch', async () => {
897+
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
898+
const fullBatch = Array.from(
899+
{ length: MEMBER_SCOPE_RENEWAL_PREFIX_BATCH },
900+
(_, index) => `source:container-${index}:`
901+
)
902+
mocks.scopes
903+
.mockResolvedValueOnce({ prefixes: fullBatch, nextCursor: 'page-2' })
904+
.mockResolvedValueOnce({ prefixes: ['source:container-last:'] })
905+
mocks.renew
906+
.mockResolvedValueOnce({ renewed: 5, finished: true })
907+
.mockResolvedValueOnce({ renewed: 1, finished: false })
908+
expect((await run()).observationsRenewed).toBe(6)
909+
expect(mocks.renew.mock.calls.map(([call]) => call.scopePrefixes.length)).toEqual([
910+
MEMBER_SCOPE_RENEWAL_PREFIX_BATCH,
911+
1,
912+
])
913+
expect(scopeRenewal()).toEqual({
914+
scopeRenewalCursor: 'page-2',
915+
scopeRenewalStartedAt: expect.any(Date),
916+
})
917+
})
918+
878919
it('does not renew again while the last renewal is recent', async () => {
879920
const run = arrange({
880921
connectorType: 'scoped_listing',
@@ -894,7 +935,7 @@ describe('member engine with a dedicated content credential', () => {
894935
contentFresh: true,
895936
scopeRenewal: { cursor: 'page-7', startedAt: new Date('2026-09-01T00:00:00Z') },
896937
})
897-
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], nextCursor: 'page-8' })
938+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'] })
898939
mocks.renew.mockResolvedValue({ renewed: 1000, finished: false })
899940
expect((await run()).observationsRenewed).toBe(1000)
900941
expect(mocks.scopes.mock.calls[0]?.[2]).toBe('page-7')
@@ -923,6 +964,16 @@ describe('member engine with a dedicated content credential', () => {
923964
})
924965
})
925966

967+
it('stops a scope listing whose cursor does not advance', async () => {
968+
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
969+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], nextCursor: 'page-2' })
970+
mocks.renew.mockResolvedValue({ renewed: 1, finished: true })
971+
expect((await run()).error).toBeUndefined()
972+
expect(mocks.scopes).toHaveBeenCalledTimes(2)
973+
expect(mocks.renew).not.toHaveBeenCalled()
974+
expect(mocks.observe).toHaveBeenCalled()
975+
})
976+
926977
it('restarts a pass whose stored cursor expired', async () => {
927978
const run = arrange({
928979
connectorType: 'scoped_listing',

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

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -76,6 +76,7 @@ import {
7676
MEMBER_FULL_RECRAWL_MINUTES,
7777
MEMBER_SCOPE_RENEW_AFTER_MS,
7878
MEMBER_SCOPE_RENEWAL_BUDGET_MS,
79+
MEMBER_SCOPE_RENEWAL_PREFIX_BATCH,
7980
MEMBER_SUSPENDED_PURGE_DAYS,
8081
MEMBER_SYNC_MAX_PAGES_PER_MEMBER,
8182
MEMBER_SYNC_SOFT_BUDGET_SECONDS,
@@ -1039,6 +1040,9 @@ async function renewMemberAccessScopes(input: {
10391040
/** A resumed pass keeps its start, so the watermark never claims more than the whole pass renewed. */
10401041
const passStartedAt = member.scopeRenewalStartedAt ?? now
10411042
let cursor = member.scopeRenewalCursor ?? undefined
1043+
/** Where the prefixes not yet renewed were read from; an unfinished pass resumes there. */
1044+
let batchCursor = cursor
1045+
let pending: string[] = []
10421046
let restartedExpiredCursor = false
10431047
let scopes = 0
10441048
let renewed = 0
@@ -1068,23 +1072,30 @@ async function renewMemberAccessScopes(input: {
10681072
)
10691073
throw error
10701074
cursor = undefined
1075+
batchCursor = undefined
1076+
pending = []
10711077
restartedExpiredCursor = true
10721078
continue
10731079
}
1080+
if (page.nextCursor && page.nextCursor === cursor)
1081+
throw new Error('Access scope pagination did not advance')
10741082
scopes += page.prefixes.length
1083+
pending.push(...page.prefixes)
1084+
cursor = page.nextCursor
1085+
if (cursor && pending.length < MEMBER_SCOPE_RENEWAL_PREFIX_BATCH) continue
10751086
const renewal = await renewMemberObservationsInScopes({
10761087
connectorId: run.connectorId,
10771088
memberId: member.id,
1078-
scopePrefixes: page.prefixes,
1089+
scopePrefixes: pending,
10791090
renewBefore,
10801091
deadlineAt,
10811092
beforeBatch: run.lease.beatIfDue,
10821093
withLease: (fn) => withMemberLease(run, fn),
10831094
})
10841095
renewed += renewal.renewed
1085-
/** An unfinished page is read again next time, from the cursor that produced it. */
10861096
if (!renewal.finished) break
1087-
cursor = page.nextCursor
1097+
pending = []
1098+
batchCursor = cursor
10881099
if (!cursor) {
10891100
finished = true
10901101
await saveProgress({
@@ -1097,7 +1108,7 @@ async function renewMemberAccessScopes(input: {
10971108
}
10981109
if (!finished)
10991110
await saveProgress({
1100-
scopeRenewalCursor: cursor ?? null,
1111+
scopeRenewalCursor: batchCursor ?? null,
11011112
scopeRenewalStartedAt: passStartedAt,
11021113
})
11031114
} catch (error) {
@@ -1626,6 +1637,7 @@ async function completeMemberSync(
16261637
docsUnchanged: result.docsUnchanged,
16271638
docsHydratedOnce: result.docsHydratedOnce,
16281639
observationsAdded: result.observationsAdded,
1640+
observationsRenewed: result.observationsRenewed,
16291641
observationsRemoved: result.observationsRemoved,
16301642
docsTombstoned: result.docsTombstoned,
16311643
docsResurrected: result.docsResurrected,
@@ -1679,6 +1691,7 @@ async function failMemberSyncLog(runId: string, result: MemberSyncResult, errorM
16791691
docsUnchanged: result.docsUnchanged,
16801692
docsHydratedOnce: result.docsHydratedOnce,
16811693
observationsAdded: result.observationsAdded,
1694+
observationsRenewed: result.observationsRenewed,
16821695
observationsRemoved: result.observationsRemoved,
16831696
docsTombstoned: result.docsTombstoned,
16841697
docsResurrected: result.docsResurrected,

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

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -731,6 +731,35 @@ describe('persistHashOnlyUpdates', () => {
731731
).toEqual({ type: 'update', existingId: 'doc-1' })
732732
})
733733

734+
it('writes refreshed source metadata alongside an equivalent content hash', async () => {
735+
const { persistHashOnlyUpdates } = await import('@/lib/knowledge/connectors/sync-persistence')
736+
queueTableRows(schemaMock.knowledgeBase, [{ id: 'kb-1' }])
737+
queueTableRows(schemaMock.knowledgeConnector, [{ id: 'connector-1' }])
738+
dbChainMockFns.returning.mockResolvedValueOnce([{ id: 'doc-1' }])
739+
const sourceModifiedAt = new Date('2026-09-08T12:00:00Z')
740+
741+
await persistHashOnlyUpdates(
742+
'kb-1',
743+
'connector-1',
744+
[
745+
{
746+
existingId: 'doc-1',
747+
externalId: 'thread-1',
748+
contentHash: 'slack-thread:v5:version:text',
749+
sourceMetadata: { sourceUrl: null, sourceModifiedAt, date1: sourceModifiedAt },
750+
},
751+
],
752+
lease
753+
)
754+
755+
expect(dbChainMockFns.set).toHaveBeenCalledWith({
756+
contentHash: 'slack-thread:v5:version:text',
757+
sourceUrl: null,
758+
sourceModifiedAt,
759+
date1: sourceModifiedAt,
760+
})
761+
})
762+
734763
it('commits live retry hashes when another document is no longer a connector target', async () => {
735764
const { persistHashOnlyUpdates } = await import('@/lib/knowledge/connectors/sync-persistence')
736765
queueTableRows(schemaMock.knowledgeBase, [{ id: 'kb-1' }])

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

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,13 @@ export const MEMBER_SCOPE_RENEW_AFTER_MS = SOURCE_ACL_MAX_AGE_MS / 2
8585
/** How much of a run one member's scope renewal may use before its listing starts. */
8686
export const MEMBER_SCOPE_RENEWAL_BUDGET_MS = 10 * 60 * 1000
8787

88+
/**
89+
* Container prefixes gathered from the source before one pass over the member's stale
90+
* observations renews them, so that pass runs once per this many containers rather than
91+
* once per source page.
92+
*/
93+
export const MEMBER_SCOPE_RENEWAL_PREFIX_BATCH = 5000
94+
8895
/** Pages applied per member before its durable feed cursor is saved for continuation. */
8996
export const MEMBER_SYNC_MAX_PAGES_PER_MEMBER = 25
9097

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

Lines changed: 28 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -435,16 +435,41 @@ export async function persistSkippedDocuments(
435435
return persisted
436436
}
437437

438+
/** Source-derived fields that can change without the indexed text changing. */
439+
export type SourceMetadataFields = Partial<DocumentTags> & {
440+
sourceUrl: string | null
441+
sourceModifiedAt: Date | null
442+
}
443+
444+
/** The source-derived fields a refreshed document carries, as a full update writes them. */
445+
export function resolveSourceMetadataFields(
446+
connectorType: string,
447+
extDoc: Pick<ExternalDocument, 'sourceUrl' | 'metadata'>,
448+
sourceConfig: Record<string, unknown>
449+
): SourceMetadataFields {
450+
return {
451+
sourceUrl: extDoc.sourceUrl ?? null,
452+
sourceModifiedAt: resolveSourceModifiedAt(extDoc.metadata),
453+
...(extDoc.metadata ? resolveTagMapping(connectorType, extDoc.metadata, sourceConfig) : {}),
454+
}
455+
}
456+
438457
/**
439458
* Persists only a new hash for existing documents, leaving indexed content and
440459
* processing state as they are: a connector-owned retry hash for a skipped
441460
* refresh, so unchanged listing metadata still re-enters hydration, or the
442-
* current hash of content that hydration found unchanged under an older one.
461+
* current hash of content that hydration found unchanged under an older one,
462+
* together with that hydration's source metadata, which can move on its own.
443463
*/
444464
export async function persistHashOnlyUpdates(
445465
knowledgeBaseId: string,
446466
connectorId: string,
447-
updates: Array<{ existingId: string; externalId: string; contentHash: string }>,
467+
updates: Array<{
468+
existingId: string
469+
externalId: string
470+
contentHash: string
471+
sourceMetadata?: SourceMetadataFields
472+
}>,
448473
lease: SyncWriteLease
449474
): Promise<string[]> {
450475
if (updates.length === 0) return []
@@ -461,7 +486,7 @@ export async function persistHashOnlyUpdates(
461486
for (const update of updates) {
462487
const persisted = await tx
463488
.update(document)
464-
.set({ contentHash: update.contentHash })
489+
.set({ contentHash: update.contentHash, ...update.sourceMetadata })
465490
.where(connectorDocumentSyncTarget(update.existingId, knowledgeBaseId, connectorId))
466491
.returning({ id: document.id })
467492
if (persisted.length === 0) {

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

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,13 @@ const mocks = vi.hoisted(() => ({
99
update: vi.fn(),
1010
triggerAvailable: vi.fn(),
1111
persistHashes: vi.fn(),
12+
sourceMetadata: vi.fn(
13+
(_connectorType: string, doc: Pick<ExternalDocument, 'sourceUrl' | 'metadata'>) => ({
14+
sourceUrl: doc.sourceUrl ?? null,
15+
sourceModifiedAt: null,
16+
date1: doc.metadata?.lastActivity,
17+
})
18+
),
1219
dispatch: vi.fn<(documents: DocumentData[]) => Promise<{ accepted: number; failed: number }>>(),
1320
}))
1421

@@ -17,6 +24,7 @@ vi.mock('@/lib/knowledge/connectors/sync-persistence', () => ({
1724
updateDocument: mocks.update,
1825
persistSkippedDocuments: vi.fn(),
1926
persistHashOnlyUpdates: mocks.persistHashes,
27+
resolveSourceMetadataFields: mocks.sourceMetadata,
2028
}))
2129
vi.mock('@/lib/knowledge/documents/service', () => ({
2230
isTriggerAvailable: mocks.triggerAvailable,
@@ -250,13 +258,15 @@ describe('processDocOps unchanged content under a new hash', () => {
250258
input.hydration.getDocument = vi.fn(async () => ({
251259
...sourceDocument('thread'),
252260
contentHash: hydratedHash,
261+
sourceUrl: 'https://source.fixture.test/thread',
262+
metadata: { lastActivity: '2026-09-08T12:00:00.000Z' },
253263
}))
254264
input.matchContentHash = (candidate, stored) =>
255265
candidate.split(':').at(-1) === stored.split(':').at(-1) ? 'equivalent' : 'stale'
256266
return input
257267
}
258268

259-
it('advances only the stored hash when the connector finds the same text', async () => {
269+
it('advances the stored hash and source metadata when the connector finds the same text', async () => {
260270
const input = refreshOf('legacy:text-a', 'version:2:text-a')
261271
await expect(processDocOps(input)).resolves.toBe(true)
262272
expect(mocks.update).not.toHaveBeenCalled()
@@ -265,7 +275,18 @@ describe('processDocOps unchanged content under a new hash', () => {
265275
expect(mocks.persistHashes).toHaveBeenCalledWith(
266276
'knowledge-base',
267277
'connector',
268-
[{ existingId: 'document', externalId: 'thread', contentHash: 'version:2:text-a' }],
278+
[
279+
{
280+
existingId: 'document',
281+
externalId: 'thread',
282+
contentHash: 'version:2:text-a',
283+
sourceMetadata: {
284+
sourceUrl: 'https://source.fixture.test/thread',
285+
sourceModifiedAt: null,
286+
date1: '2026-09-08T12:00:00.000Z',
287+
},
288+
},
289+
],
269290
input.lease
270291
)
271292
})

0 commit comments

Comments
 (0)