Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 14 additions & 2 deletions apps/sim/lib/credential-groups/slack-managed-users.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@ import {
} from '@/lib/credential-groups/slack-managed-user-scopes'
import { acquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
import type { DbOrTx } from '@/lib/db/types'
import {
listOrganizationSearchApprovals,
lockOrganizationSearchApproval,
} from '@/lib/knowledge/search/integration-policy'
import { SLACK_CUSTOM_BOT_PROVIDER_ID, SLACK_CUSTOM_BOT_SECRET_TYPE } from '@/lib/oauth/types'
import { resolveSlackAppCredentials } from '@/lib/slack-search/app-configuration'
import { requireSlackSearchAppAvailable } from '@/lib/slack-search/shared-app'
Expand Down Expand Up @@ -579,7 +583,10 @@ export async function createSlackManagedUsersAttempt(params: {
)
.limit(1)
searchApproval = {
approved: approval?.approved ?? false,
approved:
approval?.approved ??
(await listOrganizationSearchApprovals(scope.organizationId)).get('slack') ??
false,
Comment thread
waleedlatif1 marked this conversation as resolved.
updatedAt: approval?.updatedAt.getTime() ?? null,
}
if (searchApproval.approved)
Expand Down Expand Up @@ -769,6 +776,7 @@ export async function exchangeAndConfigureSlackManagedUsers(params: {

return db.transaction(async (tx) => {
if (params.attempt.organizationId) {
await lockOrganizationSearchApproval(tx, params.attempt.organizationId)
const [app] = await tx
.select()
.from(slackApp)
Expand Down Expand Up @@ -837,9 +845,13 @@ export async function exchangeAndConfigureSlackManagedUsers(params: {
)
.limit(1)
.for('share')
const approved =
approval?.approved ??
(await listOrganizationSearchApprovals(params.attempt.organizationId, tx)).get('slack') ??
Comment thread
waleedlatif1 marked this conversation as resolved.
false
Comment thread
waleedlatif1 marked this conversation as resolved.
if (
!params.attempt.searchApproval ||
(approval?.approved ?? false) !== params.attempt.searchApproval.approved ||
approved !== params.attempt.searchApproval.approved ||
(approval?.updatedAt.getTime() ?? null) !== params.attempt.searchApproval.updatedAt
)
throw new SlackManagedUsersError(
Expand Down
187 changes: 187 additions & 0 deletions apps/sim/lib/knowledge/__integration__/search-mcp-setup.integration.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,10 +50,12 @@ import {
resolveManagedOAuthToken,
} from '@/lib/credentials/managed-oauth'
import { acquireAdvisoryXactLock, tryAcquireAdvisoryXactLock } from '@/lib/db/advisory-locks'
import { deleteKnowledgeConnector } from '@/lib/knowledge/application/connectors'
import {
approveSearchIntegration,
listSearchIntegrations,
} from '@/lib/knowledge/application/search-integrations'
import { deleteKnowledgeBase, restoreKnowledgeBase } from '@/lib/knowledge/service'
import {
GITHUB_INSTALLATION_PROVIDER_ID,
type GitHubInstallationBinding,
Expand Down Expand Up @@ -332,6 +334,191 @@ describe('atomic organization live Search MCP setup', () => {
restoreSlackHttp = () => spy.mockRestore()
}

async function seedImplicitSlackApproval() {
const knowledgeBaseId = generateId()
const connectorId = generateId()
await db.insert(knowledgeBase).values({
id: knowledgeBaseId,
userId: ids.owner,
organizationId: ids.organization,
isSearchIndex: true,
name: 'Slack Search fixture',
})
await db.insert(knowledgeConnector).values({
id: connectorId,
knowledgeBaseId,
connectorType: 'slack',
status: 'active',
sourceConfig: {},
})
return connectorId
}

it.each([false, true])(
'verifies implicitly approved Search permissions unless explicitly disabled (disabled: %s)',
async (disabled) => {
const setup = await seedSlackAuthorization()
await seedImplicitSlackApproval()
if (disabled)
await approveSearchIntegration.execute({
principal: createSessionPrincipal({ userId: ids.owner, sessionId: generateId() }),
input: { organizationId: ids.organization, connectorType: 'slack', approved: false },
})
expect(await integrationStatus('slack')).toMatchObject({ approved: !disabled })
const pending = await setup.start()
const scopes = disabled
? [...SLACK_MANAGED_USER_SCOPES]
: [...new Set([...SLACK_MANAGED_USER_SCOPES, ...SLACK_SEARCH_USER_SCOPES])]
expect(new URL(pending.authorizationUrl).searchParams.get('user_scope')!.split(',')).toEqual(
expect.arrayContaining(scopes)
)
if (disabled)
expect(new URL(pending.authorizationUrl).searchParams.get('user_scope')).not.toContain(
'search:read.public'
)
provideSlackConsent(scopes)
await expect(setup.complete(pending.state)).resolves.toMatchObject({ ok: true })
const state = await snapshot()
expect(
state.groups[0].options.find((entry) => entry.id === setup.optionId)?.requiredScopes
).toEqual(expect.arrayContaining(scopes))
if (disabled)
await expect(setup.resolveToken()).resolves.toMatchObject({ accessToken: 'fixture-token' })
}
)

it.each(['added', 'removed'] as const)(
'rejects pending authorization when implicit Search approval is %s',
async (change) => {
const setup = await seedSlackAuthorization()
const connectorId = change === 'removed' ? await seedImplicitSlackApproval() : null
const pending = await setup.start()
if (connectorId)
await db
.update(knowledgeConnector)
.set({ archivedAt: new Date() })
.where(eq(knowledgeConnector.id, connectorId))
else await seedImplicitSlackApproval()
provideSlackConsent([...SLACK_MANAGED_USER_SCOPES, ...SLACK_SEARCH_USER_SCOPES])
await expect(setup.complete(pending.state)).rejects.toThrow('Search approval changed')
expect((await snapshot()).groups).toEqual(setup.before.groups)
await expect(setup.resolveToken()).resolves.toMatchObject({ accessToken: 'fixture-token' })
}
)

it('stops granting implicit Search approval when a connector is removed with documents kept', async () => {
const setup = await seedSlackAuthorization()
const connectorId = await seedImplicitSlackApproval()
await deleteKnowledgeConnector.execute({
principal: createSessionPrincipal({ userId: ids.owner, sessionId: generateId() }),
input: { connectorId, assertedOrganizationId: ids.organization, deleteDocuments: false },
})
expect(await integrationStatus('slack')).toMatchObject({ approved: false })
const pending = await setup.start()
expect(new URL(pending.authorizationUrl).searchParams.get('user_scope')).not.toContain(
'search:read.public'
)
await setup.complete(pending.state, 'access_denied')
await expect(setup.resolveToken()).resolves.toMatchObject({ accessToken: 'fixture-token' })
})

it.each(['remove connector', 'archive index', 'disable approval', 'restore index'] as const)(
'serializes the Slack consent commit with Search lifecycle changes: %s',
async (change) => {
const setup = await seedSlackAuthorization()
const connectorId = await seedImplicitSlackApproval()
const [index] = await db
.select()
.from(knowledgeBase)
.where(eq(knowledgeBase.organizationId, ids.organization))
if (change === 'restore index')
await deleteKnowledgeBase(index.id, generateId(), { allowSearchIndexDelete: true })
const pending = await setup.start()
provideSlackConsent([...SLACK_MANAGED_USER_SCOPES, ...SLACK_SEARCH_USER_SCOPES])
const probe = `slack_consent_${generateId().replace(/-/g, '')}`
const lockKey = `slack-consent-fixture:${setup.groupId}`
await db.$client.unsafe(`CREATE FUNCTION ${probe}() RETURNS trigger LANGUAGE plpgsql AS $$
BEGIN
PERFORM pg_advisory_xact_lock(hashtextextended('${lockKey}', 0));
RETURN NEW;
END $$`)
await db.$client.unsafe(`CREATE TRIGGER ${probe} BEFORE UPDATE ON credential_group
FOR EACH ROW WHEN (OLD.id = '${setup.groupId}') EXECUTE FUNCTION ${probe}()`)
const locked = createDeferred<number>()
const release = createDeferred<void>()
const blocker = db.transaction(async (tx) => {
await acquireAdvisoryXactLock(tx, 'slack_consent_fixture', lockKey)
const [connection] = await tx.execute<{ pid: number }>(sql`SELECT pg_backend_pid() AS pid`)
locked.resolve(connection.pid)
await release.promise
})
const blockerPid = await locked.promise
const callback = setup.complete(pending.state).catch((error: unknown) => error)
let mutation: Promise<unknown> | undefined
try {
let callbackPid: number | undefined
await vi.waitFor(
async () => {
const [waiting] = await db.execute<{ pid: number }>(sql`
SELECT pid FROM pg_stat_activity WHERE ${blockerPid} = ANY(pg_blocking_pids(pid))
`)
expect(waiting).toBeDefined()
callbackPid = waiting?.pid
},
{ timeout: 5_000 }
)
mutation = (
change === 'remove connector'
? deleteKnowledgeConnector.execute({
principal: createSessionPrincipal({ userId: ids.owner, sessionId: generateId() }),
input: {
connectorId,
assertedOrganizationId: ids.organization,
deleteDocuments: false,
},
})
: change === 'archive index'
? deleteKnowledgeBase(index.id, generateId(), { allowSearchIndexDelete: true })
: change === 'restore index'
? restoreKnowledgeBase(index.id, generateId())
: approveSearchIntegration.execute({
principal: createSessionPrincipal({
userId: ids.owner,
sessionId: generateId(),
}),
input: {
organizationId: ids.organization,
connectorType: 'slack',
approved: false,
},
})
).catch((error: unknown) => error)
await vi.waitFor(
async () => {
const [state] = await db.execute<{ waiting: boolean }>(sql`
SELECT EXISTS (SELECT 1 FROM pg_stat_activity
WHERE ${callbackPid!} = ANY(pg_blocking_pids(pid))) AS waiting
`)
expect(state.waiting).toBe(true)
},
{ timeout: 3_000 }
)
} finally {
release.resolve()
await blocker
const callbackResult = await callback
const mutationResult = await mutation
await db.$client.unsafe(`DROP TRIGGER ${probe} ON credential_group`)
await db.$client.unsafe(`DROP FUNCTION ${probe}()`)
expect(callbackResult).toMatchObject({ ok: true })
expect(mutationResult).not.toBeInstanceOf(Error)
}
expect(await integrationStatus('slack')).toMatchObject({
approved: change === 'restore index',
})
}
)

it.each([
{ name: 'workflow policy', scopes: SLACK_MANAGED_USER_SCOPES },
{ name: 'custom policy', scopes: ['chat:write', 'users:read', 'users:read.email'] },
Expand Down
75 changes: 38 additions & 37 deletions apps/sim/lib/knowledge/application/search-integrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,10 @@ import { resolveKnowledgeAccessAvailability } from '@/lib/knowledge/access/avail
import { defineAuthorizedKnowledgeUseCase } from '@/lib/knowledge/application/authorized-knowledge-use-case'
import { resolveKnowledgeOwnerContext } from '@/lib/knowledge/application/contexts'
import { knowledgeOperations } from '@/lib/knowledge/application/operations'
import { listOrganizationSearchApprovals } from '@/lib/knowledge/search/integration-policy'
import {
listOrganizationSearchApprovals,
lockOrganizationSearchApproval,
} from '@/lib/knowledge/search/integration-policy'
import { GITHUB_INSTALLATION_PROVIDER_ID } from '@/lib/oauth/github-installation-types'
import { refuseCapability } from '@/lib/permission-groups/capabilities'
import { isOrganizationCapabilityWithheld } from '@/lib/permission-groups/capability-assertions'
Expand Down Expand Up @@ -266,43 +269,41 @@ export const approveSearchIntegration = defineAuthorizedKnowledgeUseCase({
? await prepareSearchMcpProvider(context.organizationId, mcpProvider)
: null
let memberAccounts: { groupId: string; changed: boolean } | undefined
const changed =
policy || memberProvider || mcpProvider
? await db.transaction(async (tx) => {
if (memberProvider)
memberAccounts = await addOrganizationAccountProvider(
context.organizationId!,
requirePrincipalSubjectUserId(principal),
{
provider: memberProvider,
label: source[1].name,
...(memberProvider === 'slack'
? { requiredScopes: [...SLACK_SEARCH_USER_SCOPES] }
: {}),
},
tx
).catch((error: unknown) => {
if (error instanceof CredentialGroupProviderConfigurationError)
throw new OrchestrationError('validation', error.message)
throw error
})
if (mcpSetup)
memberAccounts = await addOrganizationSearchMcpProvider(
context.organizationId!,
requirePrincipalSubjectUserId(principal),
mcpSetup,
tx
)
if (policy)
await tx
.update(organization)
.set({
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`,
})
.where(eq(organization.id, context.organizationId!))
return saveApproval(tx)
const changed = await db.transaction(async (tx) => {
await lockOrganizationSearchApproval(tx, context.organizationId!)
if (memberProvider)
memberAccounts = await addOrganizationAccountProvider(
context.organizationId!,
requirePrincipalSubjectUserId(principal),
{
provider: memberProvider,
label: source[1].name,
...(memberProvider === 'slack'
? { requiredScopes: [...SLACK_SEARCH_USER_SCOPES] }
: {}),
},
tx
).catch((error: unknown) => {
if (error instanceof CredentialGroupProviderConfigurationError)
throw new OrchestrationError('validation', error.message)
throw error
})
if (mcpSetup)
memberAccounts = await addOrganizationSearchMcpProvider(
context.organizationId!,
requirePrincipalSubjectUserId(principal),
mcpSetup,
tx
)
if (policy)
await tx
.update(organization)
.set({
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`,
})
: await saveApproval(db)
.where(eq(organization.id, context.organizationId!))
return saveApproval(tx)
})
return {
connectorType: input.connectorType,
approved: input.approved,
Expand Down
3 changes: 3 additions & 0 deletions apps/sim/lib/knowledge/orchestration/connectors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ import {
type KnowledgeOperationContext,
type KnowledgeOrchestrationResult,
} from '@/lib/knowledge/orchestration/shared'
import { lockOrganizationSearchApproval } from '@/lib/knowledge/search/integration-policy'
import { createTagDefinition } from '@/lib/knowledge/tags/service'
import { captureServerEvent } from '@/lib/posthog/server'
import { searchSourceIdentity } from '@/lib/sim-search/source-identity'
Expand Down Expand Up @@ -469,6 +470,7 @@ export async function performCreateKnowledgeConnector(
let reused = false
try {
created = await db.transaction(async (tx) => {
if (owner.organizationId) await lockOrganizationSearchApproval(tx, owner.organizationId)
await tx.execute(sql`SELECT 1 FROM knowledge_base WHERE id = ${kb.id} FOR UPDATE`)

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