Skip to content

Commit f15fb7c

Browse files
committed
fix(knowledge): resume member scope renewal across runs
Accessible scopes are now listed page by page, and a renewal that does not finish within its budget stores its channel-listing cursor and start time, so the next run continues from there instead of re-reading the first pages and never reaching channels past them. The watermark records when the whole pass began, and an expired cursor restarts the pass once.
1 parent b50bbe8 commit f15fb7c

9 files changed

Lines changed: 223 additions & 98 deletions

File tree

‎apps/sim/connectors/slack/slack.test.ts‎

Lines changed: 21 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -708,35 +708,34 @@ describe('Slack change detection and access scopes', () => {
708708
it('reports exactly the conversations each member can read, under the listing rules', async () => {
709709
const listScopes = slackConnector.listAccessibleScopes
710710
if (!listScopes) throw new Error('Slack must report access scopes')
711+
const readAll = async (token: string, sourceConfig: Record<string, unknown>) => {
712+
const prefixes: string[] = []
713+
let cursor: string | undefined
714+
do {
715+
const page = await listScopes(token, sourceConfig, cursor)
716+
prefixes.push(...page.prefixes)
717+
cursor = page.nextCursor
718+
} while (cursor)
719+
return prefixes
720+
}
711721
pageSize = 1
712722
expect(await listScopes('alice', {})).toEqual({
713-
prefixes: [scope(GENERAL.id), scope(PRIVATE.id), scope(ARCHIVE.id)],
714-
complete: true,
723+
prefixes: [scope(GENERAL.id)],
724+
nextCursor: '1',
715725
})
716-
expect(await listScopes('bob', {})).toEqual({ prefixes: [scope(GENERAL.id)], complete: true })
726+
expect(await readAll('alice', {})).toEqual([
727+
scope(GENERAL.id),
728+
scope(PRIVATE.id),
729+
scope(ARCHIVE.id),
730+
])
731+
expect(await readAll('bob', {})).toEqual([scope(GENERAL.id)])
717732
expect(
718-
await listScopes('alice', { excludeChannels: '#general', includeArchived: 'false' })
719-
).toEqual({ prefixes: [scope(PRIVATE.id)], complete: true })
733+
await readAll('alice', { excludeChannels: '#general', includeArchived: 'false' })
734+
).toEqual([scope(PRIVATE.id)])
720735
const listed = await listAll('alice', { maxMessages: 0 })
721-
const { prefixes } = await listScopes('alice', {})
736+
const prefixes = await readAll('alice', {})
722737
expect(
723738
listed.documents.every((doc) => prefixes.some((prefix) => doc.externalId.startsWith(prefix)))
724739
).toBe(true)
725740
})
726-
727-
it('reports a conversation listing cut off at the page cap as incomplete', async () => {
728-
const listScopes = slackConnector.listAccessibleScopes
729-
if (!listScopes) throw new Error('Slack must report access scopes')
730-
replacement = (call) =>
731-
call.method === 'conversations.list'
732-
? {
733-
ok: true,
734-
channels: [GENERAL],
735-
response_metadata: { next_cursor: String(Number(call.params.get('cursor') || 0) + 1) },
736-
}
737-
: undefined
738-
const scopes = await listScopes('alice', {})
739-
expect(scopes.complete).toBe(false)
740-
expect(scopes.prefixes.length).toBeGreaterThan(0)
741-
})
742741
})

‎apps/sim/connectors/slack/slack.ts‎

Lines changed: 8 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -19,7 +19,7 @@ import {
1919
type ConnectorSourceFailureCategory,
2020
} from '@/connectors/source-error'
2121
import type {
22-
AccessibleScopes,
22+
AccessibleScopePage,
2323
ConnectorConfig,
2424
ExternalDocument,
2525
ExternalDocumentList,
@@ -43,8 +43,6 @@ const MAX_THREAD_PAGES = 200
4343
const MAX_USERNAME_CACHE_ENTRIES = 2000
4444
const MAX_CHANNEL_CACHE_ENTRIES = 2000
4545
const MAX_ROOT_VERSION_CACHE_ENTRIES = 50_000
46-
/** Conversation pages read when resolving a caller's reachable channels; each proves access on its own. */
47-
const MAX_SCOPE_CHANNEL_PAGES = 100
4846
/**
4947
* Threads with activity this recent are reread on every listing. A root reports new
5048
* replies and its own edits, but not an edit or deletion of an existing reply, which in
@@ -1008,26 +1006,19 @@ async function messageLink(
10081006
}
10091007

10101008
/**
1011-
* The conversations the caller can read, under the same inclusion rules as a listing,
1012-
* as the external-id prefix of their threads. Access is granted per conversation, so
1013-
* a conversation listed here proves access to every thread in it. Stopping at the page
1014-
* cap reports the listing incomplete, so the conversations past it stay due for renewal.
1009+
* One page of the conversations the caller can read, under the same inclusion rules as
1010+
* a listing, as the external-id prefix of their threads. Access is granted per
1011+
* conversation, so a conversation listed here proves access to every thread in it.
10151012
*/
10161013
async function listAccessibleScopes(
10171014
accessToken: string,
10181015
sourceConfig: Record<string, unknown>,
1016+
cursor?: string,
10191017
syncContext?: Record<string, unknown>
1020-
): Promise<AccessibleScopes> {
1018+
): Promise<AccessibleScopePage> {
10211019
const { teamId } = await resolveWorkspace(accessToken, syncContext)
1022-
const prefixes: string[] = []
1023-
let cursor: string | undefined
1024-
for (let page = 0; page < MAX_SCOPE_CHANNEL_PAGES; page += 1) {
1025-
const { channels, nextCursor: next } = await listChannelPage(accessToken, sourceConfig, cursor)
1026-
for (const channel of channels) prefixes.push(channelScope(teamId, channel.id))
1027-
cursor = next
1028-
if (!cursor) break
1029-
}
1030-
return { prefixes, complete: !cursor }
1020+
const { channels, nextCursor: next } = await listChannelPage(accessToken, sourceConfig, cursor)
1021+
return { prefixes: channels.map((channel) => channelScope(teamId, channel.id)), nextCursor: next }
10311022
}
10321023

10331024
export const slackConnector: ConnectorConfig = {

‎apps/sim/connectors/types.ts‎

Lines changed: 10 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -238,13 +238,13 @@ export type ExternalChange =
238238
| { kind: 'removed'; externalId: string }
239239

240240
/**
241-
* The containers a caller can still read, as external-ID prefixes. Each listed
242-
* prefix proves access on its own; `complete` is false when the source stopped
243-
* before listing every container, so the rest are still due for renewal.
241+
* One page of the containers a caller can still read, as external-ID prefixes.
242+
* Each listed prefix proves access on its own; `nextCursor` continues the
243+
* listing and is absent on its last page.
244244
*/
245-
export interface AccessibleScopes {
245+
export interface AccessibleScopePage {
246246
prefixes: string[]
247-
complete: boolean
247+
nextCursor?: string
248248
}
249249

250250
export interface ExternalChangeList {
@@ -582,14 +582,17 @@ export interface ConnectorConfig extends ConnectorMeta {
582582
* granted to whole containers (a Slack channel) rather than item by item. A
583583
* members-mode crawl renews the caller's observations under these prefixes
584584
* without relisting each item, so access stays fresh while a large listing is
585-
* still in progress, and lapses for containers the caller has lost. Only
585+
* still in progress, and lapses for containers the caller has lost. Paged, so
586+
* a renewal can resume where an earlier run stopped; an expired cursor is
587+
* recognised by {@link ConnectorConfig.isListingCursorInvalidError}. Only
586588
* meaningful alongside {@link ConnectorMeta.permissionScopedListing}.
587589
*/
588590
listAccessibleScopes?: (
589591
accessToken: string,
590592
sourceConfig: Record<string, unknown>,
593+
cursor?: string,
591594
syncContext?: Record<string, unknown>
592-
) => Promise<AccessibleScopes>
595+
) => Promise<AccessibleScopePage>
593596

594597
/**
595598
* Opens the external directory whose groups this connector's mirrored ACLs

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

Lines changed: 79 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -98,6 +98,8 @@ vi.mock('@/connectors/registry.server', () => ({
9898
listDocuments: mocks.list,
9999
getDocument: mocks.get,
100100
listAccessibleScopes: mocks.scopes,
101+
isListingCursorInvalidError: (error: unknown) =>
102+
error instanceof Error && error.message === 'cursor expired',
101103
},
102104
drive: {
103105
id: 'drive',
@@ -166,6 +168,7 @@ function arrange(
166168
organizationId?: string
167169
memberCheckpoint?: Record<string, unknown>
168170
scopeRenewedAt?: Date
171+
scopeRenewal?: { cursor: string; startedAt: Date }
169172
} = {}
170173
) {
171174
const connector = {
@@ -213,6 +216,8 @@ function arrange(
213216
: {
214217
...member,
215218
scopeRenewedAt: options.scopeRenewedAt ?? null,
219+
scopeRenewalCursor: options.scopeRenewal?.cursor ?? null,
220+
scopeRenewalStartedAt: options.scopeRenewal?.startedAt ?? null,
216221
...(options.memberCheckpoint ? { listingCheckpoint: options.memberCheckpoint } : {}),
217222
}
218223
const claims = options.members && !options.noDueMembers ? [[memberRow], []] : [[]]
@@ -358,6 +363,8 @@ describe('member engine with a dedicated content credential', () => {
358363
memberSyncedThrough: null,
359364
lastCompleteListingAt: null,
360365
scopeRenewedAt: null,
366+
scopeRenewalCursor: null,
367+
scopeRenewalStartedAt: null,
361368
listingCheckpoint: { kind: 'membership', cursor: null, removeMember: false },
362369
changeCursor: null,
363370
nextAttemptAt: expect.any(Date),
@@ -706,6 +713,8 @@ describe('member engine with a dedicated content credential', () => {
706713
memberSyncedThrough: null,
707714
lastCompleteListingAt: null,
708715
scopeRenewedAt: null,
716+
scopeRenewalCursor: null,
717+
scopeRenewalStartedAt: null,
709718
listingCheckpoint: { kind: 'membership', cursor: null, removeMember: false },
710719
})
711720
)
@@ -826,11 +835,13 @@ describe('member engine with a dedicated content credential', () => {
826835
})
827836

828837
const scopeRenewal = () =>
829-
dbChainMockFns.set.mock.calls.find(([value]) => 'scopeRenewedAt' in value)?.[0]
838+
dbChainMockFns.set.mock.calls.find(
839+
([value]) => 'scopeRenewedAt' in value || 'scopeRenewalCursor' in value
840+
)?.[0]
830841

831842
it('renews member access by scope before listing and records when it finished', async () => {
832843
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
833-
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], complete: true })
844+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'] })
834845
mocks.renew.mockResolvedValue({ renewed: 3, finished: true })
835846
const result = await run()
836847
expect(result.error).toBeUndefined()
@@ -842,7 +853,26 @@ describe('member engine with a dedicated content credential', () => {
842853
expect(mocks.renew.mock.invocationCallOrder[0]).toBeLessThan(
843854
mocks.list.mock.invocationCallOrder[memberListing]
844855
)
845-
expect(scopeRenewal()).toEqual({ scopeRenewedAt: expect.any(Date) })
856+
expect(scopeRenewal()).toEqual({
857+
scopeRenewedAt: expect.any(Date),
858+
scopeRenewalCursor: null,
859+
scopeRenewalStartedAt: null,
860+
})
861+
})
862+
863+
it('renews every page of scopes before recording the renewal', async () => {
864+
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
865+
mocks.scopes
866+
.mockResolvedValueOnce({ prefixes: ['source:container-a:'], nextCursor: 'page-2' })
867+
.mockResolvedValueOnce({ prefixes: ['source:container-b:'] })
868+
mocks.renew.mockResolvedValue({ renewed: 2, finished: true })
869+
expect((await run()).observationsRenewed).toBe(4)
870+
expect(mocks.scopes.mock.calls.map((call) => call[2])).toEqual([undefined, 'page-2'])
871+
expect(mocks.renew.mock.calls.map(([call]) => call.scopePrefixes)).toEqual([
872+
['source:container-a:'],
873+
['source:container-b:'],
874+
])
875+
expect(scopeRenewal()).toMatchObject({ scopeRenewedAt: expect.any(Date) })
846876
})
847877

848878
it('does not renew again while the last renewal is recent', async () => {
@@ -857,23 +887,56 @@ describe('member engine with a dedicated content credential', () => {
857887
expect(mocks.renew).not.toHaveBeenCalled()
858888
})
859889

860-
it('leaves an unfinished renewal due for the next run', async () => {
861-
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
862-
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], complete: true })
890+
it('stores the cursor of an unfinished page so the next run reads it again', async () => {
891+
const run = arrange({
892+
connectorType: 'scoped_listing',
893+
members: true,
894+
contentFresh: true,
895+
scopeRenewal: { cursor: 'page-7', startedAt: new Date('2026-09-01T00:00:00Z') },
896+
})
897+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], nextCursor: 'page-8' })
863898
mocks.renew.mockResolvedValue({ renewed: 1000, finished: false })
864899
expect((await run()).observationsRenewed).toBe(1000)
865-
expect(scopeRenewal()).toBeUndefined()
900+
expect(mocks.scopes.mock.calls[0]?.[2]).toBe('page-7')
901+
expect(scopeRenewal()).toEqual({
902+
scopeRenewalCursor: 'page-7',
903+
scopeRenewalStartedAt: new Date('2026-09-01T00:00:00Z'),
904+
})
866905
})
867906

868-
it('leaves renewal due when the source stopped before listing every scope', async () => {
869-
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
870-
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], complete: false })
871-
mocks.renew.mockResolvedValue({ renewed: 3, finished: true })
872-
expect((await run()).observationsRenewed).toBe(3)
873-
expect(mocks.renew).toHaveBeenCalledWith(
874-
expect.objectContaining({ scopePrefixes: ['source:container-a:'] })
875-
)
876-
expect(scopeRenewal()).toBeUndefined()
907+
it('resumes a stored pass and records its original start once every page is renewed', async () => {
908+
const startedAt = new Date('2026-09-01T00:00:00Z')
909+
const run = arrange({
910+
connectorType: 'scoped_listing',
911+
members: true,
912+
contentFresh: true,
913+
scopeRenewal: { cursor: 'page-7', startedAt },
914+
})
915+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-z:'] })
916+
mocks.renew.mockResolvedValue({ renewed: 1, finished: true })
917+
await run()
918+
expect(mocks.scopes.mock.calls.map((call) => call[2])).toEqual(['page-7'])
919+
expect(scopeRenewal()).toEqual({
920+
scopeRenewedAt: startedAt,
921+
scopeRenewalCursor: null,
922+
scopeRenewalStartedAt: null,
923+
})
924+
})
925+
926+
it('restarts a pass whose stored cursor expired', async () => {
927+
const run = arrange({
928+
connectorType: 'scoped_listing',
929+
members: true,
930+
contentFresh: true,
931+
scopeRenewal: { cursor: 'page-7', startedAt: new Date('2026-09-01T00:00:00Z') },
932+
})
933+
mocks.scopes
934+
.mockRejectedValueOnce(new Error('cursor expired'))
935+
.mockResolvedValueOnce({ prefixes: ['source:container-a:'] })
936+
mocks.renew.mockResolvedValue({ renewed: 1, finished: true })
937+
await run()
938+
expect(mocks.scopes.mock.calls.map((call) => call[2])).toEqual(['page-7', undefined])
939+
expect(scopeRenewal()).toMatchObject({ scopeRenewalCursor: null })
877940
})
878941

879942
it('still lists the member when their access scopes cannot be read', async () => {

0 commit comments

Comments
 (0)