Skip to content

Commit c78bd1c

Browse files
committed
fix(knowledge): keep Slack searchable while its member crawl runs
Members-mode search shows a document only while its member observation is younger than a day, but a full Slack listing re-fetched every thread and so took far longer than a day on a large workspace, leaving most of Slack invisible. Access is now renewed per channel the member can still read, and listings re-read only threads whose root changed, are active, or are due in a rolling 28-day refresh.
1 parent 5eece6e commit c78bd1c

23 files changed

Lines changed: 28942 additions & 93 deletions

‎.github/workflows/test-build.yml‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -246,6 +246,8 @@ jobs:
246246
lib/knowledge/__integration__/stored-document-recovery.integration.ts
247247
lib/knowledge/__integration__/connector-partition-work.integration.ts
248248
lib/knowledge/__integration__/listing-continuation.integration.ts
249+
lib/knowledge/__integration__/member-scope-renewal.integration.ts
250+
lib/knowledge/__integration__/slack-empty-threads.integration.ts
249251
lib/knowledge/__integration__/kb-block-search.integration.ts
250252
lib/core/outbox/service.integration.ts
251253
lib/knowledge/__integration__/connector-upload.integration.ts

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -31,7 +31,7 @@ vi.mock('@/lib/knowledge/documents/service', () => ({
3131
vi.mock('@/lib/knowledge/connectors/sync-persistence', () => ({
3232
addDocument: vi.fn(),
3333
persistSkippedDocuments: vi.fn(),
34-
persistSkippedRetryHashes: vi.fn(),
34+
persistHashOnlyUpdates: vi.fn(),
3535
updateDocument: vi.fn(),
3636
}))
3737

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

Lines changed: 138 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
*/
44
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
55
import { slackConnectorMeta } from '@/connectors/slack/meta'
6-
import { slackConnector } from '@/connectors/slack/slack'
6+
import { slackConnector, slackRollingRefreshDay } from '@/connectors/slack/slack'
77
import type { ExternalDocument } from '@/connectors/types'
88
import { CONNECTOR_TEXT_DOCUMENT_MAX_BYTES } from '@/connectors/utils'
99

@@ -166,8 +166,19 @@ beforeEach(() => {
166166

167167
afterEach(() => {
168168
vi.unstubAllGlobals()
169+
vi.useRealTimers()
169170
})
170171

172+
const DAY = 24 * 60 * 60 * 1000
173+
174+
/** Freezes the clock on a day that is not this thread's rolling reread, so it lists as quiet. */
175+
function freezeOnQuietDay(externalId: string): void {
176+
const base = Date.UTC(2026, 0, 1)
177+
const offset = (slackRollingRefreshDay(externalId) + 1 - (Math.floor(base / DAY) % 28) + 28) % 28
178+
vi.useFakeTimers({ toFake: ['Date'] })
179+
vi.setSystemTime(base + offset * DAY)
180+
}
181+
171182
async function listAll(token: string, config: Record<string, unknown> = {}, run = 'run-1') {
172183
const context: Record<string, unknown> = { syncRunId: run }
173184
const documents: ExternalDocument[] = []
@@ -275,12 +286,13 @@ describe('Slack thread indexing through provider APIs', () => {
275286
expect((await listAll('alice', { channel: GENERAL.id })).documents).toEqual([])
276287
})
277288

278-
it('retires legacy channel ids and marks every thread for reply refresh on the next run', async () => {
289+
it('retires legacy channel ids and keeps a quiet thread under one hash across listings', async () => {
290+
freezeOnQuietDay(id(GENERAL.id))
279291
const first = await listAll('alice', { channel: GENERAL.id }, 'run-1')
280292
const next = await listAll('alice', { channel: GENERAL.id }, 'run-2')
281293
expect(first.documents.map((doc) => doc.externalId)).not.toContain(GENERAL.id)
282294
expect(first.documents[0].externalId).toBe(next.documents[0].externalId)
283-
expect(first.documents[0].contentHash).not.toBe(next.documents[0].contentHash)
295+
expect(first.documents[0].contentHash).toBe(next.documents[0].contentHash)
284296
expect(await slackConnector.getDocument('alice', {}, GENERAL.id)).toBeNull()
285297
})
286298

@@ -590,3 +602,126 @@ describe('Slack incomplete and unsafe provider responses', () => {
590602
)
591603
})
592604
})
605+
606+
describe('Slack change detection and access scopes', () => {
607+
const match = (candidate: string, stored: string) =>
608+
slackConnector.matchContentHash?.(candidate, stored)
609+
const scope = (channel: string) => `slack:v4:${TEAM}:${channel}:`
610+
611+
beforeEach(() => {
612+
freezeOnQuietDay(id(GENERAL.id))
613+
})
614+
615+
it('rereads each quiet thread on exactly one day of every 28', async () => {
616+
const start = Date.now()
617+
const rereadDays: number[] = []
618+
for (let day = 0; day < 28; day += 1) {
619+
vi.setSystemTime(start + day * DAY)
620+
const listed = await listAll('alice', { channel: GENERAL.id }, `run-${day}`)
621+
if (listed.documents[0].contentHash.includes(':refresh:')) rereadDays.push(day)
622+
}
623+
expect(rereadDays).toHaveLength(1)
624+
})
625+
626+
it('rereads every thread on a full resync without changing the hash of unchanged text', async () => {
627+
const quiet = await listAll('alice', { channel: GENERAL.id }, 'run-1')
628+
const stored = await slackConnector.getDocument('alice', {}, id(GENERAL.id), quiet.context)
629+
const context: Record<string, unknown> = { syncRunId: 'full-run', fullSync: true }
630+
const full = await slackConnector.listDocuments(
631+
'alice',
632+
{ channel: GENERAL.id },
633+
undefined,
634+
context
635+
)
636+
expect(match(full.documents[0].contentHash, stored?.contentHash ?? '')).toBe('stale')
637+
const reread = await slackConnector.getDocument('alice', {}, id(GENERAL.id), context)
638+
expect(reread?.contentHash).toBe(stored?.contentHash)
639+
})
640+
641+
it('does not reread a quiet thread until its root reports a change', async () => {
642+
const first = await listAll('alice', { channel: GENERAL.id }, 'run-1')
643+
const stored = await slackConnector.getDocument('alice', {}, id(GENERAL.id), first.context)
644+
const next = await listAll('alice', { channel: GENERAL.id }, 'run-2')
645+
expect(match(next.documents[0].contentHash, stored?.contentHash ?? '')).toBe('current')
646+
channels[0].messages = [{ ...root(), reply_count: 2, latest_reply: SECOND }]
647+
const replied = await listAll('alice', { channel: GENERAL.id }, 'run-3')
648+
expect(match(replied.documents[0].contentHash, stored?.contentHash ?? '')).toBe('stale')
649+
channels[0].messages = [{ ...root(), edited: { ts: SECOND } }]
650+
const edited = await listAll('alice', { channel: GENERAL.id }, 'run-4')
651+
expect(match(edited.documents[0].contentHash, stored?.contentHash ?? '')).toBe('stale')
652+
})
653+
654+
it('rereads an active thread every listing and changes its hash only when the text changed', async () => {
655+
const recent = `${Math.floor(Date.now() / 1000) - 3600}.000100`
656+
channels[0].messages = [{ ...root(), latest_reply: recent }]
657+
const first = await listAll('alice', { channel: GENERAL.id }, 'run-1')
658+
const stored = await slackConnector.getDocument('alice', {}, id(GENERAL.id), first.context)
659+
const next = await listAll('alice', { channel: GENERAL.id }, 'run-2')
660+
expect(next.documents[0].contentHash).not.toBe(first.documents[0].contentHash)
661+
expect(match(next.documents[0].contentHash, stored?.contentHash ?? '')).toBe('stale')
662+
const reread = await slackConnector.getDocument('alice', {}, id(GENERAL.id), next.context)
663+
expect(reread?.contentHash).toBe(stored?.contentHash)
664+
channels[0].replies[ROOT] = [root(), reply('Use the corrected queue')]
665+
const edited = await slackConnector.getDocument('alice', {}, id(GENERAL.id), next.context)
666+
expect(match(edited?.contentHash ?? '', stored?.contentHash ?? '')).toBe('stale')
667+
})
668+
669+
it('leaves a verified empty quiet thread settled until its root changes', async () => {
670+
channels[0].messages = [{ ...root(''), thread_ts: ROOT }]
671+
channels[0].replies[ROOT] = [{ ...root(''), thread_ts: ROOT }]
672+
const listed = await listAll('alice', { channel: GENERAL.id }, 'run-1')
673+
expect(listed.documents[0].skippedRetryPolicy).toBe('source-change')
674+
const skipped = await slackConnector.getDocument('alice', {}, id(GENERAL.id), listed.context)
675+
expect(skipped?.skippedReason).toBeDefined()
676+
const next = await listAll('alice', { channel: GENERAL.id }, 'run-2')
677+
expect(match(next.documents[0].contentHash, skipped?.contentHash ?? '')).toBe('current')
678+
channels[0].messages = [{ ...root(''), thread_ts: ROOT, reply_count: 2, latest_reply: SECOND }]
679+
const replied = await listAll('alice', { channel: GENERAL.id }, 'run-3')
680+
expect(match(replied.documents[0].contentHash, skipped?.contentHash ?? '')).toBe('stale')
681+
})
682+
683+
it('carries a thread indexed under the text-only hash forward without reindexing it', async () => {
684+
const hydrated = await slackConnector.getDocument('alice', {}, id(GENERAL.id))
685+
const text = hydrated?.contentHash.split(':').at(-1)
686+
expect(match(hydrated?.contentHash ?? '', `slack-content:v4:${text}`)).toBe('equivalent')
687+
expect(match(hydrated?.contentHash ?? '', `slack-content:v4:${'0'.repeat(64)}`)).toBe('stale')
688+
const listed = await listAll('alice', { channel: GENERAL.id })
689+
expect(match(listed.documents[0].contentHash, `slack-content:v4:${text}`)).toBe('stale')
690+
})
691+
692+
it('links messages from the workspace address and reads each conversation once per run', async () => {
693+
replacement = (call) =>
694+
call.method === 'auth.test'
695+
? { ok: true, team_id: TEAM, url: 'https://acme.slack.com/' }
696+
: undefined
697+
channels[0].messages = [root(), { ...root('Second thread'), ts: SECOND }]
698+
channels[0].replies[SECOND] = [{ ...root('Second thread'), ts: SECOND }]
699+
const context: Record<string, unknown> = {}
700+
const first = await slackConnector.getDocument('alice', {}, id(GENERAL.id), context)
701+
await slackConnector.getDocument('alice', {}, id(GENERAL.id, SECOND), context)
702+
expect(first?.sourceUrl).toBe('https://acme.slack.com/archives/C0GENERAL/p1700000100000100')
703+
expect(calls.some((call) => call.method === 'chat.getPermalink')).toBe(false)
704+
expect(calls.filter((call) => call.method === 'conversations.info')).toHaveLength(1)
705+
expect(calls.filter((call) => call.method === 'auth.test')).toHaveLength(1)
706+
})
707+
708+
it('reports exactly the conversations each member can read, under the listing rules', async () => {
709+
const listScopes = slackConnector.listAccessibleScopes
710+
if (!listScopes) throw new Error('Slack must report access scopes')
711+
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)])
718+
expect(
719+
await listScopes('alice', { excludeChannels: '#general', includeArchived: 'false' })
720+
).toEqual([scope(PRIVATE.id)])
721+
const listed = await listAll('alice', { maxMessages: 0 })
722+
const scopes = await listScopes('alice', {})
723+
expect(
724+
listed.documents.every((doc) => scopes.some((prefix) => doc.externalId.startsWith(prefix)))
725+
).toBe(true)
726+
})
727+
})

0 commit comments

Comments
 (0)