Skip to content

Commit dd9c165

Browse files
fix(slack-search): finalize failed streams and explain query limits
1 parent aa1155e commit dd9c165

10 files changed

Lines changed: 369 additions & 26 deletions

File tree

‎apps/sim/lib/knowledge/application/slack-search/assistant.test.ts‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,8 @@ const m = vi.hoisted(() => ({
1717
start: vi.fn(),
1818
finish: vi.fn(),
1919
finishWithError: vi.fn(),
20+
terminate: vi.fn(),
21+
streamOptions: vi.fn(),
2022
payload: vi.fn(),
2123
outcome: vi.fn(),
2224
memberAuthorization: vi.fn(),
@@ -99,9 +101,13 @@ vi.mock('@/lib/copilot/request/session/abort', () => ({
99101
}))
100102
vi.mock('@/lib/slack-search/assistant-stream', () => ({
101103
SlackSearchAssistantStream: class {
104+
constructor(options: unknown) {
105+
m.streamOptions(options)
106+
}
102107
start = m.start
103108
finish = m.finish
104109
finishWithError = m.finishWithError
110+
terminateAfterFailure = m.terminate
105111
onEvent = vi.fn()
106112
assertHealthy = vi.fn()
107113
},
@@ -308,9 +314,24 @@ describe('organization Assistant from Slack', () => {
308314
m.finish.mockRejectedValueOnce(new Error('ambiguous send'))
309315
await expect(run()).rejects.toThrow('ambiguous send')
310316
expect(m.run).toHaveBeenCalledOnce()
317+
expect(m.terminate).toHaveBeenCalledOnce()
311318
expect(m.updateRun).toHaveBeenCalledWith('run1', 'error')
312319
expect(m.release).toHaveBeenCalledOnce()
313320
})
321+
it('rechecks cleanup authority with a fresh signal after aborting execution', async () => {
322+
m.run.mockRejectedValueOnce(new Error('Assistant disconnected'))
323+
m.terminate.mockImplementationOnce(async () => {
324+
const options = m.streamOptions.mock.calls[0][0]
325+
expect(options.controller.signal.aborted).toBe(true)
326+
await options.beforeCleanup(new AbortController().signal)
327+
})
328+
await expect(run()).rejects.toThrow('Assistant disconnected')
329+
expect(m.terminate).toHaveBeenCalledOnce()
330+
expect(m.memberAuthorization).toHaveBeenCalled()
331+
expect(m.terminate.mock.invocationCallOrder[0]).toBeLessThan(
332+
m.release.mock.invocationCallOrder[0]
333+
)
334+
})
314335
it('closes a healthy Slack stream when the Assistant reports failure', async () => {
315336
m.run.mockResolvedValueOnce({ success: false, content: '', contentBlocks: [], toolCalls: [] })
316337
await expect(run()).rejects.toThrow('Organization Assistant did not complete')
@@ -328,9 +349,10 @@ describe('organization Assistant from Slack', () => {
328349
})
329350
it('marks private history as stopped only when the turn was cancelled by native Stop', async () => {
330351
m.run.mockRejectedValueOnce(new Error('Stream stopped'))
331-
m.stopped.mockResolvedValueOnce(true)
352+
m.stopped.mockResolvedValue(true)
332353
await expect(run()).rejects.toThrow('Stream stopped')
333354
expect(m.stoppedMessage).toHaveBeenCalledOnce()
334355
expect(m.finishWithError).not.toHaveBeenCalled()
356+
expect(m.terminate).not.toHaveBeenCalled()
335357
})
336358
})

‎apps/sim/lib/knowledge/application/slack-search/assistant.ts‎

Lines changed: 20 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -140,9 +140,10 @@ export async function runSlackSearchAssistant(
140140
let failed = true
141141
let failure: Error | undefined
142142
let registry: ResolvedSecretTraceRegistry | undefined
143+
let stream: SlackSearchAssistantStream | undefined
143144
let checking = false
144-
const checkAccess = async () => {
145-
controller.signal.throwIfAborted()
145+
const checkAccess = async (signal: AbortSignal = controller.signal) => {
146+
signal.throwIfAborted()
146147
await requireSlackSearchTurnLease(turnId, leaseId)
147148
if (!(await authorizeSlackSearchInstallation(principal, job)))
148149
throw new OrchestrationError('forbidden', 'Slack Search is disabled')
@@ -223,15 +224,17 @@ export async function runSlackSearchAssistant(
223224
})
224225
if (!run) throw new Error('Could not persist Assistant execution')
225226
runId = run.id
226-
const stream = new SlackSearchAssistantStream({
227+
const responseStream = new SlackSearchAssistantStream({
227228
token: secret.botToken,
228229
channel: job.message.channelId,
229230
threadTs: job.message.threadTs ?? job.message.messageTs,
230231
slackUserId: job.message.userId,
231232
controller,
232233
registry: environmentContext.resolvedSecretTraceRegistry,
233234
beforeDelivery: checkAccess,
235+
beforeCleanup: checkAccess,
234236
})
237+
stream = responseStream
235238
const payload = await buildCopilotRequestPayload(
236239
{
237240
message: job.message.query,
@@ -248,7 +251,7 @@ export async function runSlackSearchAssistant(
248251
actorUserId: userId,
249252
organizationId: installation.organizationId,
250253
})
251-
await stream.start()
254+
await responseStream.start()
252255
result = await runHeadlessCopilotLifecycle(payload, {
253256
userId,
254257
organizationId: installation.organizationId,
@@ -264,28 +267,36 @@ export async function runSlackSearchAssistant(
264267
autoExecuteTools: true,
265268
onEvent: async (event) => {
266269
try {
267-
await stream.onEvent(event)
270+
await responseStream.onEvent(event)
268271
} catch (error) {
269272
controller.abort(error)
270273
throw error
271274
}
272275
},
273276
})
274-
stream.assertHealthy()
277+
responseStream.assertHealthy()
275278
controller.signal.throwIfAborted()
276279
if (!result.success) {
277-
await stream.finishWithError()
280+
await responseStream.finishWithError()
278281
throw new Error('Organization Assistant did not complete')
279282
}
280283
await checkAccess()
281-
await stream.finish(result)
284+
await responseStream.finish(result)
282285
failed = false
283286
await recordSlackSearchOutcome(installation, 'success')
284287
}
285288
} catch (error) {
286289
failure = toError(error)
287290
controller.abort(error)
291+
const failedStream = stream
288292
const outcomes = await Promise.allSettled([
293+
...(failedStream
294+
? [
295+
wasSlackSearchTurnStopped(turnId, leaseId).then((stopped) =>
296+
stopped ? undefined : failedStream.terminateAfterFailure()
297+
),
298+
]
299+
: []),
289300
recordSlackSearchOutcome(installation, 'assistant_or_delivery_failed'),
290301
...(runId ? [updateRunStatus(runId, 'error')] : []),
291302
])
@@ -295,7 +306,7 @@ export async function runSlackSearchAssistant(
295306
if (errors.length)
296307
failure = new AggregateError(
297308
[failure, ...errors],
298-
'Slack turn and outcome persistence failed'
309+
'Slack turn cleanup or outcome persistence failed'
299310
)
300311
} finally {
301312
await titleTask
Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,100 @@
1+
/** @vitest-environment node */
2+
import { dbChainMockFns, queueTableRows, resetDbChainMock, schemaMock } from '@sim/testing'
3+
import { beforeEach, describe, expect, it, vi } from 'vitest'
4+
5+
const mocks = vi.hoisted(() => ({ authorize: vi.fn(), openDm: vi.fn(), post: vi.fn() }))
6+
vi.mock('@/lib/knowledge/application/slack-search/authorization', () => ({
7+
authorizeSlackSearchInstallation: mocks.authorize,
8+
}))
9+
vi.mock('@/lib/internal/slack/client', () => ({
10+
openSlackDm: mocks.openDm,
11+
postSlackMessage: mocks.post,
12+
slackString: (value: Record<string, unknown>, key: string) =>
13+
typeof value[key] === 'string' ? value[key] : undefined,
14+
}))
15+
16+
import { routeSlackSearchMentionToDm } from '@/lib/knowledge/application/slack-search/mention'
17+
import { SLACK_SEARCH_QUERY_TOO_LONG } from '@/lib/slack-search/constants'
18+
import { slackSearchJobSchema } from '@/lib/slack-search/types'
19+
20+
const principal = {
21+
kind: 'slack_installation',
22+
credentialId: 'c1',
23+
credentialVersion: 'v1',
24+
appId: 'A1',
25+
teamId: 'T1',
26+
eventId: 'Ev1',
27+
receivedAt: new Date(),
28+
} as const
29+
30+
beforeEach(() => {
31+
vi.clearAllMocks()
32+
resetDbChainMock()
33+
mocks.authorize.mockResolvedValue({ secret: { botToken: 'test-token' } })
34+
mocks.openDm.mockResolvedValue('D1')
35+
mocks.post.mockResolvedValue({
36+
status: 200,
37+
data: { ok: true, channel: 'D1', ts: '1800000000.2' },
38+
})
39+
})
40+
41+
describe('private Slack mention roots', () => {
42+
it.each([false, true])(
43+
'creates a nonempty private root for queryTooLong=%s',
44+
async (queryTooLong) => {
45+
const job = slackSearchJobSchema.parse({
46+
installationId: 'i1',
47+
revision: 'r1',
48+
credentialId: 'c1',
49+
credentialVersion: 'v1',
50+
receivedAt: principal.receivedAt.getTime(),
51+
message: {
52+
appId: 'A1',
53+
teamId: 'T1',
54+
eventId: 'Ev1',
55+
channelId: 'C1',
56+
userId: 'U1',
57+
messageTs: '1800000000.1',
58+
origin: { channelId: 'C1', threadTs: '1800000000.1', messageTs: '1800000000.1' },
59+
query: queryTooLong ? '' : 'Find the release notes',
60+
queryTooLong,
61+
},
62+
})
63+
queueTableRows(schemaMock.slackSearchInstallation, [
64+
{ enabled: true, revision: 'r1', credentialVersion: 'v1' },
65+
])
66+
queueTableRows(schemaMock.slackSearchTurn, [
67+
{
68+
status: 'running',
69+
leaseId: 'lease1',
70+
leaseExpiresAt: new Date(Date.now() + 60_000),
71+
payload: job,
72+
},
73+
])
74+
const routed = await routeSlackSearchMentionToDm(principal, {
75+
job,
76+
turnId: 'turn1',
77+
leaseId: 'lease1',
78+
signal: new AbortController().signal,
79+
})
80+
expect(mocks.post).toHaveBeenCalledWith(
81+
'test-token',
82+
expect.objectContaining({
83+
channel: 'D1',
84+
blocks: [
85+
{
86+
type: 'section',
87+
text: {
88+
type: 'plain_text',
89+
text: queryTooLong ? SLACK_SEARCH_QUERY_TOO_LONG : job.message.query,
90+
},
91+
},
92+
],
93+
}),
94+
expect.any(AbortSignal)
95+
)
96+
expect(routed.message).toMatchObject({ channelId: 'D1', threadTs: '1800000000.2' })
97+
expect(dbChainMockFns.set).toHaveBeenCalledWith(expect.objectContaining({ payload: routed }))
98+
}
99+
)
100+
})

‎apps/sim/lib/knowledge/application/slack-search/mention.ts‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import { and, eq } from 'drizzle-orm'
55
import { OrchestrationError } from '@/lib/core/orchestration/types'
66
import { openSlackDm, postSlackMessage, slackString } from '@/lib/internal/slack/client'
77
import { authorizeSlackSearchInstallation } from '@/lib/knowledge/application/slack-search/authorization'
8+
import { SLACK_SEARCH_QUERY_TOO_LONG } from '@/lib/slack-search/constants'
89
import { slackSearchConversationKey } from '@/lib/slack-search/conversation'
910
import { type SlackSearchJob, slackSearchJobSchema } from '@/lib/slack-search/types'
1011

@@ -72,8 +73,18 @@ export async function routeSlackSearchMentionToDm(
7273
context.secret.botToken,
7374
{
7475
channel: channelId,
75-
text: 'Your question for Sim Search',
76-
blocks: [{ type: 'section', text: { type: 'plain_text', text: job.message.query } }],
76+
text: job.message.queryTooLong
77+
? SLACK_SEARCH_QUERY_TOO_LONG
78+
: 'Your question for Sim Search',
79+
blocks: [
80+
{
81+
type: 'section',
82+
text: {
83+
type: 'plain_text',
84+
text: job.message.queryTooLong ? SLACK_SEARCH_QUERY_TOO_LONG : job.message.query,
85+
},
86+
},
87+
],
7788
unfurl_links: false,
7889
unfurl_media: false,
7990
},

0 commit comments

Comments
 (0)