Skip to content

Commit d3f0d26

Browse files
authored
fix(slack): honor implicit Search approval during authorization (#8548)
* fix(slack): honor implicit Search approval during authorization * fix(slack): serialize Search approval changes with consent
1 parent 320843e commit d3f0d26

7 files changed

Lines changed: 268 additions & 44 deletions

File tree

‎apps/sim/lib/credential-groups/slack-managed-users.ts‎

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,10 @@ import {
3333
} from '@/lib/credential-groups/slack-managed-user-scopes'
3434
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
3535
import type { DbOrTx } from '@/lib/db/types'
36+
import {
37+
listOrganizationSearchApprovals,
38+
lockOrganizationSearchApproval,
39+
} from '@/lib/knowledge/search/integration-policy'
3640
import { SLACK_CUSTOM_BOT_PROVIDER_ID, SLACK_CUSTOM_BOT_SECRET_TYPE } from '@/lib/oauth/types'
3741
import { resolveSlackAppCredentials } from '@/lib/slack-search/app-configuration'
3842
import { requireSlackSearchAppAvailable } from '@/lib/slack-search/shared-app'
@@ -579,7 +583,10 @@ export async function createSlackManagedUsersAttempt(params: {
579583
)
580584
.limit(1)
581585
searchApproval = {
582-
approved: approval?.approved ?? false,
586+
approved:
587+
approval?.approved ??
588+
(await listOrganizationSearchApprovals(scope.organizationId)).get('slack') ??
589+
false,
583590
updatedAt: approval?.updatedAt.getTime() ?? null,
584591
}
585592
if (searchApproval.approved)
@@ -769,6 +776,7 @@ export async function exchangeAndConfigureSlackManagedUsers(params: {
769776

770777
return db.transaction(async (tx) => {
771778
if (params.attempt.organizationId) {
779+
await lockOrganizationSearchApproval(tx, params.attempt.organizationId)
772780
const [app] = await tx
773781
.select()
774782
.from(slackApp)
@@ -837,9 +845,13 @@ export async function exchangeAndConfigureSlackManagedUsers(params: {
837845
)
838846
.limit(1)
839847
.for('share')
848+
const approved =
849+
approval?.approved ??
850+
(await listOrganizationSearchApprovals(params.attempt.organizationId, tx)).get('slack') ??
851+
false
840852
if (
841853
!params.attempt.searchApproval ||
842-
(approval?.approved ?? false) !== params.attempt.searchApproval.approved ||
854+
approved !== params.attempt.searchApproval.approved ||
843855
(approval?.updatedAt.getTime() ?? null) !== params.attempt.searchApproval.updatedAt
844856
)
845857
throw new SlackManagedUsersError(

‎apps/sim/lib/knowledge/__integration__/search-mcp-setup.integration.ts‎

Lines changed: 187 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,10 +50,12 @@ import {
5050
resolveManagedOAuthToken,
5151
} from '@/lib/credentials/managed-oauth'
5252
import { acquireAdvisoryXactLock, tryAcquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
53+
import { deleteKnowledgeConnector } from '@/lib/knowledge/application/connectors'
5354
import {
5455
approveSearchIntegration,
5556
listSearchIntegrations,
5657
} from '@/lib/knowledge/application/search-integrations'
58+
import { deleteKnowledgeBase, restoreKnowledgeBase } from '@/lib/knowledge/service'
5759
import {
5860
GITHUB_INSTALLATION_PROVIDER_ID,
5961
type GitHubInstallationBinding,
@@ -332,6 +334,191 @@ describe('atomic organization live Search MCP setup', () => {
332334
restoreSlackHttp = () => spy.mockRestore()
333335
}
334336

337+
async function seedImplicitSlackApproval() {
338+
const knowledgeBaseId = generateId()
339+
const connectorId = generateId()
340+
await db.insert(knowledgeBase).values({
341+
id: knowledgeBaseId,
342+
userId: ids.owner,
343+
organizationId: ids.organization,
344+
isSearchIndex: true,
345+
name: 'Slack Search fixture',
346+
})
347+
await db.insert(knowledgeConnector).values({
348+
id: connectorId,
349+
knowledgeBaseId,
350+
connectorType: 'slack',
351+
status: 'active',
352+
sourceConfig: {},
353+
})
354+
return connectorId
355+
}
356+
357+
it.each([false, true])(
358+
'verifies implicitly approved Search permissions unless explicitly disabled (disabled: %s)',
359+
async (disabled) => {
360+
const setup = await seedSlackAuthorization()
361+
await seedImplicitSlackApproval()
362+
if (disabled)
363+
await approveSearchIntegration.execute({
364+
principal: createSessionPrincipal({ userId: ids.owner, sessionId: generateId() }),
365+
input: { organizationId: ids.organization, connectorType: 'slack', approved: false },
366+
})
367+
expect(await integrationStatus('slack')).toMatchObject({ approved: !disabled })
368+
const pending = await setup.start()
369+
const scopes = disabled
370+
? [...SLACK_MANAGED_USER_SCOPES]
371+
: [...new Set([...SLACK_MANAGED_USER_SCOPES, ...SLACK_SEARCH_USER_SCOPES])]
372+
expect(new URL(pending.authorizationUrl).searchParams.get('user_scope')!.split(',')).toEqual(
373+
expect.arrayContaining(scopes)
374+
)
375+
if (disabled)
376+
expect(new URL(pending.authorizationUrl).searchParams.get('user_scope')).not.toContain(
377+
'search:read.public'
378+
)
379+
provideSlackConsent(scopes)
380+
await expect(setup.complete(pending.state)).resolves.toMatchObject({ ok: true })
381+
const state = await snapshot()
382+
expect(
383+
state.groups[0].options.find((entry) => entry.id === setup.optionId)?.requiredScopes
384+
).toEqual(expect.arrayContaining(scopes))
385+
if (disabled)
386+
await expect(setup.resolveToken()).resolves.toMatchObject({ accessToken: 'fixture-token' })
387+
}
388+
)
389+
390+
it.each(['added', 'removed'] as const)(
391+
'rejects pending authorization when implicit Search approval is %s',
392+
async (change) => {
393+
const setup = await seedSlackAuthorization()
394+
const connectorId = change === 'removed' ? await seedImplicitSlackApproval() : null
395+
const pending = await setup.start()
396+
if (connectorId)
397+
await db
398+
.update(knowledgeConnector)
399+
.set({ archivedAt: new Date() })
400+
.where(eq(knowledgeConnector.id, connectorId))
401+
else await seedImplicitSlackApproval()
402+
provideSlackConsent([...SLACK_MANAGED_USER_SCOPES, ...SLACK_SEARCH_USER_SCOPES])
403+
await expect(setup.complete(pending.state)).rejects.toThrow('Search approval changed')
404+
expect((await snapshot()).groups).toEqual(setup.before.groups)
405+
await expect(setup.resolveToken()).resolves.toMatchObject({ accessToken: 'fixture-token' })
406+
}
407+
)
408+
409+
it('stops granting implicit Search approval when a connector is removed with documents kept', async () => {
410+
const setup = await seedSlackAuthorization()
411+
const connectorId = await seedImplicitSlackApproval()
412+
await deleteKnowledgeConnector.execute({
413+
principal: createSessionPrincipal({ userId: ids.owner, sessionId: generateId() }),
414+
input: { connectorId, assertedOrganizationId: ids.organization, deleteDocuments: false },
415+
})
416+
expect(await integrationStatus('slack')).toMatchObject({ approved: false })
417+
const pending = await setup.start()
418+
expect(new URL(pending.authorizationUrl).searchParams.get('user_scope')).not.toContain(
419+
'search:read.public'
420+
)
421+
await setup.complete(pending.state, 'access_denied')
422+
await expect(setup.resolveToken()).resolves.toMatchObject({ accessToken: 'fixture-token' })
423+
})
424+
425+
it.each(['remove connector', 'archive index', 'disable approval', 'restore index'] as const)(
426+
'serializes the Slack consent commit with Search lifecycle changes: %s',
427+
async (change) => {
428+
const setup = await seedSlackAuthorization()
429+
const connectorId = await seedImplicitSlackApproval()
430+
const [index] = await db
431+
.select()
432+
.from(knowledgeBase)
433+
.where(eq(knowledgeBase.organizationId, ids.organization))
434+
if (change === 'restore index')
435+
await deleteKnowledgeBase(index.id, generateId(), { allowSearchIndexDelete: true })
436+
const pending = await setup.start()
437+
provideSlackConsent([...SLACK_MANAGED_USER_SCOPES, ...SLACK_SEARCH_USER_SCOPES])
438+
const probe = `slack_consent_${generateId().replace(/-/g, '')}`
439+
const lockKey = `slack-consent-fixture:${setup.groupId}`
440+
await db.$client.unsafe(`CREATE FUNCTION ${probe}() RETURNS trigger LANGUAGE plpgsql AS $$
441+
BEGIN
442+
PERFORM pg_advisory_xact_lock(hashtextextended('${lockKey}', 0));
443+
RETURN NEW;
444+
END $$`)
445+
await db.$client.unsafe(`CREATE TRIGGER ${probe} BEFORE UPDATE ON credential_group
446+
FOR EACH ROW WHEN (OLD.id = '${setup.groupId}') EXECUTE FUNCTION ${probe}()`)
447+
const locked = createDeferred<number>()
448+
const release = createDeferred<void>()
449+
const blocker = db.transaction(async (tx) => {
450+
await acquireAdvisoryXactLock(tx, 'slack_consent_fixture', lockKey)
451+
const [connection] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`)
452+
locked.resolve(connection.pid)
453+
await release.promise
454+
})
455+
const blockerPid = await locked.promise
456+
const callback = setup.complete(pending.state).catch((error: unknown) => error)
457+
let mutation: Promise<unknown> | undefined
458+
try {
459+
let callbackPid: number | undefined
460+
await vi.waitFor(
461+
async () => {
462+
const [waiting] = await db.execute<{ pid: number }>(sql`
463+
SELECT pid FROM pg_stat_activity WHERE ${blockerPid} = ANY(pg_blocking_pids(pid))
464+
`)
465+
expect(waiting).toBeDefined()
466+
callbackPid = waiting?.pid
467+
},
468+
{ timeout: 5_000 }
469+
)
470+
mutation = (
471+
change === 'remove connector'
472+
? deleteKnowledgeConnector.execute({
473+
principal: createSessionPrincipal({ userId: ids.owner, sessionId: generateId() }),
474+
input: {
475+
connectorId,
476+
assertedOrganizationId: ids.organization,
477+
deleteDocuments: false,
478+
},
479+
})
480+
: change === 'archive index'
481+
? deleteKnowledgeBase(index.id, generateId(), { allowSearchIndexDelete: true })
482+
: change === 'restore index'
483+
? restoreKnowledgeBase(index.id, generateId())
484+
: approveSearchIntegration.execute({
485+
principal: createSessionPrincipal({
486+
userId: ids.owner,
487+
sessionId: generateId(),
488+
}),
489+
input: {
490+
organizationId: ids.organization,
491+
connectorType: 'slack',
492+
approved: false,
493+
},
494+
})
495+
).catch((error: unknown) => error)
496+
await vi.waitFor(
497+
async () => {
498+
const [state] = await db.execute<{ waiting: boolean }>(sql`
499+
SELECT EXISTS (SELECT 1 FROM pg_stat_activity
500+
WHERE ${callbackPid!} = ANY(pg_blocking_pids(pid))) AS waiting
501+
`)
502+
expect(state.waiting).toBe(true)
503+
},
504+
{ timeout: 3_000 }
505+
)
506+
} finally {
507+
release.resolve()
508+
await blocker
509+
const callbackResult = await callback
510+
const mutationResult = await mutation
511+
await db.$client.unsafe(`DROP TRIGGER ${probe} ON credential_group`)
512+
await db.$client.unsafe(`DROP FUNCTION ${probe}()`)
513+
expect(callbackResult).toMatchObject({ ok: true })
514+
expect(mutationResult).not.toBeInstanceOf(Error)
515+
}
516+
expect(await integrationStatus('slack')).toMatchObject({
517+
approved: change === 'restore index',
518+
})
519+
}
520+
)
521+
335522
it.each([
336523
{ name: 'workflow policy', scopes: SLACK_MANAGED_USER_SCOPES },
337524
{ name: 'custom policy', scopes: ['chat:write', 'users:read', 'users:read.email'] },

‎apps/sim/lib/knowledge/application/search-integrations.ts‎

Lines changed: 38 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,10 @@ import { resolveKnowledgeAccessAvailability } from '@/lib/knowledge/access/avail
1818
import { defineAuthorizedKnowledgeUseCase } from '@/lib/knowledge/application/authorized-knowledge-use-case'
1919
import { resolveKnowledgeOwnerContext } from '@/lib/knowledge/application/contexts'
2020
import { knowledgeOperations } from '@/lib/knowledge/application/operations'
21-
import { listOrganizationSearchApprovals } from '@/lib/knowledge/search/integration-policy'
21+
import {
22+
listOrganizationSearchApprovals,
23+
lockOrganizationSearchApproval,
24+
} from '@/lib/knowledge/search/integration-policy'
2225
import { GITHUB_INSTALLATION_PROVIDER_ID } from '@/lib/oauth/github-installation-types'
2326
import { refuseCapability } from '@/lib/permission-groups/capabilities'
2427
import { isOrganizationCapabilityWithheld } from '@/lib/permission-groups/capability-assertions'
@@ -266,43 +269,41 @@ export const approveSearchIntegration = defineAuthorizedKnowledgeUseCase({
266269
? await prepareSearchMcpProvider(context.organizationId, mcpProvider)
267270
: null
268271
let memberAccounts: { groupId: string; changed: boolean } | undefined
269-
const changed =
270-
policy || memberProvider || mcpProvider
271-
? await db.transaction(async (tx) => {
272-
if (memberProvider)
273-
memberAccounts = await addOrganizationAccountProvider(
274-
context.organizationId!,
275-
requirePrincipalSubjectUserId(principal),
276-
{
277-
provider: memberProvider,
278-
label: source[1].name,
279-
...(memberProvider === 'slack'
280-
? { requiredScopes: [...SLACK_SEARCH_USER_SCOPES] }
281-
: {}),
282-
},
283-
tx
284-
).catch((error: unknown) => {
285-
if (error instanceof CredentialGroupProviderConfigurationError)
286-
throw new OrchestrationError('validation', error.message)
287-
throw error
288-
})
289-
if (mcpSetup)
290-
memberAccounts = await addOrganizationSearchMcpProvider(
291-
context.organizationId!,
292-
requirePrincipalSubjectUserId(principal),
293-
mcpSetup,
294-
tx
295-
)
296-
if (policy)
297-
await tx
298-
.update(organization)
299-
.set({
300-
metadata: sql`jsonb_set(COALESCE(${organization.metadata}::jsonb, '{}'::jsonb), '{liveSearchPolicies}', COALESCE(${organization.metadata}::jsonb->'liveSearchPolicies', '{}'::jsonb) || jsonb_build_object(${input.connectorType}::text, ${JSON.stringify(policy)}::jsonb))::json`,
301-
})
302-
.where(eq(organization.id, context.organizationId!))
303-
return saveApproval(tx)
272+
const changed = await db.transaction(async (tx) => {
273+
await lockOrganizationSearchApproval(tx, context.organizationId!)
274+
if (memberProvider)
275+
memberAccounts = await addOrganizationAccountProvider(
276+
context.organizationId!,
277+
requirePrincipalSubjectUserId(principal),
278+
{
279+
provider: memberProvider,
280+
label: source[1].name,
281+
...(memberProvider === 'slack'
282+
? { requiredScopes: [...SLACK_SEARCH_USER_SCOPES] }
283+
: {}),
284+
},
285+
tx
286+
).catch((error: unknown) => {
287+
if (error instanceof CredentialGroupProviderConfigurationError)
288+
throw new OrchestrationError('validation', error.message)
289+
throw error
290+
})
291+
if (mcpSetup)
292+
memberAccounts = await addOrganizationSearchMcpProvider(
293+
context.organizationId!,
294+
requirePrincipalSubjectUserId(principal),
295+
mcpSetup,
296+
tx
297+
)
298+
if (policy)
299+
await tx
300+
.update(organization)
301+
.set({
302+
metadata: sql`jsonb_set(COALESCE(${organization.metadata}::jsonb, '{}'::jsonb), '{liveSearchPolicies}', COALESCE(${organization.metadata}::jsonb->'liveSearchPolicies', '{}'::jsonb) || jsonb_build_object(${input.connectorType}::text, ${JSON.stringify(policy)}::jsonb))::json`,
304303
})
305-
: await saveApproval(db)
304+
.where(eq(organization.id, context.organizationId!))
305+
return saveApproval(tx)
306+
})
306307
return {
307308
connectorType: input.connectorType,
308309
approved: input.approved,

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,7 @@ import {
6767
type KnowledgeOperationContext,
6868
type KnowledgeOrchestrationResult,
6969
} from '@/lib/knowledge/orchestration/shared'
70+
import { lockOrganizationSearchApproval } from '@/lib/knowledge/search/integration-policy'
7071
import { createTagDefinition } from '@/lib/knowledge/tags/service'
7172
import { captureServerEvent } from '@/lib/posthog/server'
7273
import { searchSourceIdentity } from '@/lib/sim-search/source-identity'
@@ -469,6 +470,7 @@ export async function performCreateKnowledgeConnector(
469470
let reused = false
470471
try {
471472
created = await db.transaction(async (tx) => {
473+
if (owner.organizationId) await lockOrganizationSearchApproval(tx, owner.organizationId)
472474
await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${kb.id} FOR UPDATE`)
473475

474476
const activeKb = await tx
@@ -1220,6 +1222,7 @@ export async function performDeleteKnowledgeConnector(
12201222
docCount = await db.transaction(async (tx) => {
12211223
await tx.execute(sql`SET LOCAL lock_timeout = '5s'`)
12221224
await tx.execute(sql`SET LOCAL statement_timeout = '10s'`)
1225+
if (owner.organizationId) await lockOrganizationSearchApproval(tx, owner.organizationId)
12231226
/** Match source writes and document deletion: parent KB, connector, then storage ledgers. */
12241227
const [lockedOwner] = await tx
12251228
.select({

0 commit comments

Comments
 (0)