Skip to content

Commit cef0b6d

Browse files
committed
fix(knowledge): keep scope renewal due when a channel listing is cut off
A Slack conversation listing that stops at its page cap now reports itself incomplete, and the member's renewal watermark only advances once every reachable channel has been listed and renewed.
1 parent 1932033 commit cef0b6d

5 files changed

Lines changed: 65 additions & 22 deletions

File tree

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

Lines changed: 24 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -709,19 +709,34 @@ describe('Slack change detection and access scopes', () => {
709709
const listScopes = slackConnector.listAccessibleScopes
710710
if (!listScopes) throw new Error('Slack must report access scopes')
711711
pageSize = 1
712-
expect(await listScopes('alice', {})).toEqual([
713-
scope(GENERAL.id),
714-
scope(PRIVATE.id),
715-
scope(ARCHIVE.id),
716-
])
717-
expect(await listScopes('bob', {})).toEqual([scope(GENERAL.id)])
712+
expect(await listScopes('alice', {})).toEqual({
713+
prefixes: [scope(GENERAL.id), scope(PRIVATE.id), scope(ARCHIVE.id)],
714+
complete: true,
715+
})
716+
expect(await listScopes('bob', {})).toEqual({ prefixes: [scope(GENERAL.id)], complete: true })
718717
expect(
719718
await listScopes('alice', { excludeChannels: '#general', includeArchived: 'false' })
720-
).toEqual([scope(PRIVATE.id)])
719+
).toEqual({ prefixes: [scope(PRIVATE.id)], complete: true })
721720
const listed = await listAll('alice', { maxMessages: 0 })
722-
const scopes = await listScopes('alice', {})
721+
const { prefixes } = await listScopes('alice', {})
723722
expect(
724-
listed.documents.every((doc) => scopes.some((prefix) => doc.externalId.startsWith(prefix)))
723+
listed.documents.every((doc) => prefixes.some((prefix) => doc.externalId.startsWith(prefix)))
725724
).toBe(true)
726725
})
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+
})
727742
})

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

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,12 @@ import {
1818
ConnectorSourceError,
1919
type ConnectorSourceFailureCategory,
2020
} from '@/connectors/source-error'
21-
import type { ConnectorConfig, ExternalDocument, ExternalDocumentList } from '@/connectors/types'
21+
import type {
22+
AccessibleScopes,
23+
ConnectorConfig,
24+
ExternalDocument,
25+
ExternalDocumentList,
26+
} from '@/connectors/types'
2227
import {
2328
BoundedLines,
2429
CONNECTOR_TEXT_DOCUMENT_MAX_BYTES,
@@ -1005,23 +1010,24 @@ async function messageLink(
10051010
/**
10061011
* The conversations the caller can read, under the same inclusion rules as a listing,
10071012
* as the external-id prefix of their threads. Access is granted per conversation, so
1008-
* a conversation listed here proves access to every thread in it.
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.
10091015
*/
10101016
async function listAccessibleScopes(
10111017
accessToken: string,
10121018
sourceConfig: Record<string, unknown>,
10131019
syncContext?: Record<string, unknown>
1014-
): Promise<string[]> {
1020+
): Promise<AccessibleScopes> {
10151021
const { teamId } = await resolveWorkspace(accessToken, syncContext)
1016-
const scopes: string[] = []
1022+
const prefixes: string[] = []
10171023
let cursor: string | undefined
10181024
for (let page = 0; page < MAX_SCOPE_CHANNEL_PAGES; page += 1) {
10191025
const { channels, nextCursor: next } = await listChannelPage(accessToken, sourceConfig, cursor)
1020-
for (const channel of channels) scopes.push(channelScope(teamId, channel.id))
1026+
for (const channel of channels) prefixes.push(channelScope(teamId, channel.id))
10211027
cursor = next
10221028
if (!cursor) break
10231029
}
1024-
return scopes
1030+
return { prefixes, complete: !cursor }
10251031
}
10261032

10271033
export const slackConnector: ConnectorConfig = {

‎apps/sim/connectors/types.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -237,6 +237,16 @@ export type ExternalChange =
237237
| { kind: 'upsert'; externalId: string; document: ExternalDocument }
238238
| { kind: 'removed'; externalId: string }
239239

240+
/**
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.
244+
*/
245+
export interface AccessibleScopes {
246+
prefixes: string[]
247+
complete: boolean
248+
}
249+
240250
export interface ExternalChangeList {
241251
changes: ExternalChange[]
242252
/**
@@ -579,7 +589,7 @@ export interface ConnectorConfig extends ConnectorMeta {
579589
accessToken: string,
580590
sourceConfig: Record<string, unknown>,
581591
syncContext?: Record<string, unknown>
582-
) => Promise<string[]>
592+
) => Promise<AccessibleScopes>
583593

584594
/**
585595
* Opens the external directory whose groups this connector's mirrored ACLs

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

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -828,7 +828,7 @@ describe('member engine with a dedicated content credential', () => {
828828

829829
it('renews member access by scope before listing and records when it finished', async () => {
830830
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
831-
mocks.scopes.mockResolvedValue(['source:container-a:'])
831+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], complete: true })
832832
mocks.renew.mockResolvedValue({ renewed: 3, finished: true })
833833
const result = await run()
834834
expect(result.error).toBeUndefined()
@@ -857,12 +857,23 @@ describe('member engine with a dedicated content credential', () => {
857857

858858
it('leaves an unfinished renewal due for the next run', async () => {
859859
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
860-
mocks.scopes.mockResolvedValue(['source:container-a:'])
860+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], complete: true })
861861
mocks.renew.mockResolvedValue({ renewed: 1000, finished: false })
862862
expect((await run()).observationsRenewed).toBe(1000)
863863
expect(scopeRenewal()).toBeUndefined()
864864
})
865865

866+
it('leaves renewal due when the source stopped before listing every scope', async () => {
867+
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
868+
mocks.scopes.mockResolvedValue({ prefixes: ['source:container-a:'], complete: false })
869+
mocks.renew.mockResolvedValue({ renewed: 3, finished: true })
870+
expect((await run()).observationsRenewed).toBe(3)
871+
expect(mocks.renew).toHaveBeenCalledWith(
872+
expect.objectContaining({ scopePrefixes: ['source:container-a:'] })
873+
)
874+
expect(scopeRenewal()).toBeUndefined()
875+
})
876+
866877
it('still lists the member when their access scopes cannot be read', async () => {
867878
const run = arrange({ connectorType: 'scoped_listing', members: true, contentFresh: true })
868879
mocks.scopes.mockRejectedValue(new Error('provider unavailable'))

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

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1027,22 +1027,22 @@ async function renewMemberAccessScopes(input: {
10271027
/** Observations under lost containers stay stale, so only the member's own watermark says renewal is due. */
10281028
if (member.scopeRenewedAt && member.scopeRenewedAt > renewBefore) return
10291029
try {
1030-
const scopePrefixes = await connectorConfig.listAccessibleScopes(
1030+
const scopes = await connectorConfig.listAccessibleScopes(
10311031
await input.tokens.get(member.id),
10321032
input.sourceConfig,
10331033
input.syncContext
10341034
)
10351035
const renewal = await renewMemberObservationsInScopes({
10361036
connectorId: run.connectorId,
10371037
memberId: member.id,
1038-
scopePrefixes,
1038+
scopePrefixes: scopes.prefixes,
10391039
renewBefore,
10401040
deadlineAt: Math.min(run.deadlineAt, Date.now() + MEMBER_SCOPE_RENEWAL_BUDGET_MS),
10411041
beforeBatch: run.lease.beatIfDue,
10421042
withLease: (fn) => withMemberLease(run, fn),
10431043
})
10441044
run.result.observationsRenewed += renewal.renewed
1045-
if (renewal.finished) {
1045+
if (renewal.finished && scopes.complete) {
10461046
await withMemberLease(run, (tx) =>
10471047
tx
10481048
.update(knowledgeConnectorMember)
@@ -1053,7 +1053,8 @@ async function renewMemberAccessScopes(input: {
10531053
logger.info('Renewed member observations by access scope', {
10541054
connectorId: run.connectorId,
10551055
memberId: member.id,
1056-
scopes: scopePrefixes.length,
1056+
scopes: scopes.prefixes.length,
1057+
scopesComplete: scopes.complete,
10571058
renewed: renewal.renewed,
10581059
finished: renewal.finished,
10591060
})

0 commit comments

Comments
 (0)