Skip to content

Commit 409b869

Browse files
committed
fix(slack): deliver Sim Chat text and tool progress reliably
1 parent c7f5b6e commit 409b869

19 files changed

Lines changed: 1011 additions & 252 deletions

File tree

‎apps/docs/content/docs/integrations/slack.mdx‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1970,7 +1970,7 @@ Trigger from Slack events, interactions, and slash commands
19701970
| `streamTaskTitle` | string | No | Optional status Slack shows while each selected response is being produced. Leave empty to use Running. |
19711971
| `streamTaskDisplayMode` | string | No | Choose how Slack displays thinking and tool progress. |
19721972
| `streamIncludeThinking` | boolean | No | Show agent thinking as Slack task updates while the response is generated. |
1973-
| `streamIncludeToolCalls` | boolean | No | Show tool execution lifecycle as Slack task updates. |
1973+
| `streamIncludeToolCalls` | boolean | No | Show Agent and Sim Chat tool execution progress as Slack task updates. |
19741974
| `emoji` | string | No | Comma-separated emoji names to restrict to. Leave empty to match any emoji. |
19751975
| `nameContains` | string | No | Only fire when the created channel name contains this text. |
19761976
| `interactionFilter` | string | No | Comma-separated action_ids \(buttons/selects\) or callback_ids \(modals\) to restrict to. Leave empty to fire on any interaction. |

‎apps/sim/app/api/mothership/execute/route.test.ts‎

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -391,19 +391,18 @@ describe('mothership private trace provenance transport', () => {
391391
.map((line) => JSON.parse(line))
392392
expect(response.status).toBe(200)
393393
expect(events.filter((event) => event.type === 'error')).toEqual([])
394-
expect(events.filter((event) => event.type === 'agent_event')).toEqual([
395-
{ type: 'agent_event', event: { type: 'thinking_delta', text: 'Considering' } },
396-
{ type: 'agent_event', event: { type: 'turn_end', turn: 'intermediate' } },
397-
{ type: 'agent_event', event: { type: 'tool_call_start', id: 'tool-1', name: 'Lookup' } },
398-
{
399-
type: 'agent_event',
400-
event: { type: 'tool_call_end', id: 'tool-1', name: 'Lookup', status: 'success' },
401-
},
402-
{ type: 'agent_event', event: { type: 'turn_end', turn: 'final' } },
403-
])
394+
expect(events.filter((event) => event.type === 'agent_event')).toEqual(
395+
[
396+
{ type: 'thinking_delta', text: 'Considering' },
397+
{ type: 'turn_end', turn: 'intermediate' },
398+
{ type: 'tool_call_start', id: 'tool-1', name: 'Lookup' },
399+
{ type: 'tool_call_end', id: 'tool-1', name: 'Lookup', status: 'success' },
400+
{ type: 'turn_end', turn: 'final' },
401+
].map((event) => ({ type: 'agent_event', v: 1, event }))
402+
)
404403
expect(events.filter((event) => event.type === 'chunk')).toEqual([
405-
{ type: 'chunk', content: 'I will check.', turn: 'pending' },
406-
{ type: 'chunk', content: 'Answer', turn: 'pending' },
404+
{ type: 'chunk', v: 1, content: 'I will check.', turn: 'pending' },
405+
{ type: 'chunk', v: 1, content: 'Answer', turn: 'pending' },
407406
])
408407
expect(text).not.toContain('private-args')
409408
expect(text).not.toContain('private-result')

‎apps/sim/app/api/mothership/execute/route.ts‎

Lines changed: 18 additions & 60 deletions
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,10 @@ import { createLogger } from '@sim/logger'
33
import { getErrorMessage, toError } from '@sim/utils/errors'
44
import { generateId } from '@sim/utils/id'
55
import { type NextRequest, NextResponse } from 'next/server'
6-
import { mothershipExecuteContract } from '@/lib/api/contracts/mothership-chats'
6+
import {
7+
type MothershipExecuteStreamEvent,
8+
mothershipExecuteContract,
9+
} from '@/lib/api/contracts/mothership-chats'
710
import { parseRequest } from '@/lib/api/server'
811
import { checkInternalAuth } from '@/lib/auth/hybrid'
912
import { verifyInternalDelegationToken } from '@/lib/auth/internal'
@@ -27,11 +30,8 @@ import {
2730
createCopilotEnvironmentContext,
2831
} from '@/lib/mothership/environment-context'
2932
import { IntegrationCatalogMcpExecution } from '@/lib/mothership/generated/integration-catalog'
30-
import {
31-
MothershipStreamV1EventType,
32-
MothershipStreamV1TextChannel,
33-
} from '@/lib/mothership/generated/mothership-stream-v1'
3433
import { PROTOCOL_VERSION } from '@/lib/mothership/generated/protocol'
34+
import { ExecuteEventProjection } from '@/lib/mothership/request/lifecycle/execute-events'
3535
import { runHeadlessCopilotLifecycle } from '@/lib/mothership/request/lifecycle/headless'
3636
import { requestExplicitStreamAbort } from '@/lib/mothership/request/session/explicit-abort'
3737
import type { StreamEvent } from '@/lib/mothership/request/types'
@@ -391,15 +391,19 @@ export const POST = withRouteHandler(async (req: NextRequest) => {
391391

392392
const stream = new ReadableStream<Uint8Array>({
393393
start(controller) {
394-
let lastForwardedSeq = -1
395-
const startedTools = new Set<string>()
396-
const endedTools = new Set<string>()
397-
const send = (event: unknown) => {
394+
const send = (event: MothershipExecuteStreamEvent) => {
398395
if (!cancelled) {
399396
controller.enqueue(encodeNdjson(event))
400397
}
401398
}
402399

400+
const projection = new ExecuteEventProjection((event) => {
401+
/** Keep text frames readable by workers deployed before agent-event versioning. */
402+
if (event.type === 'text_delta')
403+
send({ type: 'chunk', v: 1, content: event.text, turn: event.turn })
404+
else send({ type: 'agent_event', v: 1, event })
405+
})
406+
403407
// Flush response headers promptly and keep long headless runs from
404408
// looking idle to worker/proxy HTTP stacks.
405409
send({ type: 'heartbeat', timestamp: new Date().toISOString() })
@@ -409,59 +413,11 @@ export const POST = withRouteHandler(async (req: NextRequest) => {
409413

410414
void (async () => {
411415
try {
412-
const result = await runLifecycle(async (event) => {
413-
/** Reconnect replays share the same monotone sequence across event types. */
414-
if (typeof event.seq === 'number') {
415-
if (event.seq <= lastForwardedSeq) return
416-
lastForwardedSeq = event.seq
417-
}
418-
if (event.type === MothershipStreamV1EventType.text && event.payload.text) {
419-
if (event.payload.channel === MothershipStreamV1TextChannel.assistant) {
420-
if (event.scope?.lane === 'subagent') return
421-
send({ type: 'chunk', content: event.payload.text, turn: 'pending' })
422-
} else if (event.payload.channel === MothershipStreamV1TextChannel.thinking) {
423-
send({
424-
type: 'agent_event',
425-
event: { type: 'thinking_delta', text: event.payload.text },
426-
})
427-
}
428-
} else if (event.type === MothershipStreamV1EventType.tool) {
429-
const tool = event.payload
430-
if (!('phase' in tool)) return
431-
if (tool.phase === 'call' && !startedTools.has(tool.toolCallId)) {
432-
startedTools.add(tool.toolCallId)
433-
if (event.scope?.lane !== 'subagent' && !tool.replay) {
434-
send({
435-
type: 'agent_event',
436-
event: { type: 'turn_end', turn: 'intermediate' },
437-
})
438-
}
439-
send({
440-
type: 'agent_event',
441-
event: { type: 'tool_call_start', id: tool.toolCallId, name: tool.toolName },
442-
})
443-
} else if (tool.phase === 'result' && !endedTools.has(tool.toolCallId)) {
444-
endedTools.add(tool.toolCallId)
445-
send({
446-
type: 'agent_event',
447-
event: {
448-
type: 'tool_call_end',
449-
id: tool.toolCallId,
450-
name: tool.toolName,
451-
status:
452-
tool.status === 'cancelled'
453-
? 'cancelled'
454-
: tool.success
455-
? 'success'
456-
: 'error',
457-
},
458-
})
459-
}
460-
}
461-
})
416+
const result = await runLifecycle(async (event) => projection.accept(event))
462417
allowExplicitAbort = false
463418

464419
if (lifecycleAbortController.signal.aborted) {
420+
projection.finish('cancelled')
465421
send(
466422
withPrivateProvenance(
467423
{ type: 'error', error: 'Sim execution aborted' },
@@ -473,6 +429,7 @@ export const POST = withRouteHandler(async (req: NextRequest) => {
473429
}
474430

475431
if (!result.success) {
432+
projection.finish('error')
476433
logger.error(
477434
messageId
478435
? `Mothership execute failed [messageId:${messageId}]`
@@ -499,7 +456,7 @@ export const POST = withRouteHandler(async (req: NextRequest) => {
499456
return
500457
}
501458

502-
send({ type: 'agent_event', event: { type: 'turn_end', turn: 'final' } })
459+
projection.finish('success')
503460
send({
504461
type: 'final',
505462
data: withPrivateProvenance(
@@ -509,6 +466,7 @@ export const POST = withRouteHandler(async (req: NextRequest) => {
509466
),
510467
})
511468
} catch (error) {
469+
projection.finish(lifecycleAbortController.signal.aborted ? 'cancelled' : 'error')
512470
if (
513471
lifecycleAbortController.signal.aborted ||
514472
req.signal.aborted ||

‎apps/sim/background/webhook-execution.test.ts‎

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,13 @@ vi.mock('@/lib/webhooks/env-resolver', () => ({
7575
resolveWebhookRecordProviderConfig: mockResolveWebhookRecordProviderConfig,
7676
}))
7777

78+
const { mockCreateSlackStreamController } = vi.hoisted(() => ({
79+
mockCreateSlackStreamController: vi.fn(),
80+
}))
81+
vi.mock('@/lib/webhooks/slack-execution-stream', () => ({
82+
SlackExecutionStreamController: { create: mockCreateSlackStreamController },
83+
}))
84+
7885
vi.mock('@/lib/workflows/executor/execution-core', () => ({
7986
executeWorkflowCore: mockExecuteWorkflowCore,
8087
wasExecutionFinalizedByCore: mockWasExecutionFinalizedByCore,
@@ -324,6 +331,59 @@ describe('executeWebhookJob fault vs error handling', () => {
324331
dbChainMockFns.limit.mockResolvedValue([{ id: 'webhook-1' }])
325332
})
326333

334+
it.each([true, false])(
335+
'enables the shared Agent/Sim Chat event stream and honors tool visibility: %s',
336+
async (includeToolCalls) => {
337+
const config = {
338+
enabled: true,
339+
outputConfigs: [{ blockId: 'mship', path: 'content' }],
340+
includeToolCalls,
341+
includeThinking: false,
342+
taskDisplayMode: 'timeline',
343+
taskTitle: 'Running',
344+
}
345+
dbChainMockFns.limit.mockResolvedValue([
346+
{
347+
id: 'webhook-1',
348+
providerConfig: { credentialId: 'bot-1', streamResponseConfig: config },
349+
},
350+
])
351+
const calls: string[] = []
352+
const controller = {
353+
selectedOutputs: ['mship_content'],
354+
callbacks: { onStream: vi.fn() },
355+
finalize: vi.fn(async () => {
356+
calls.push('delivery')
357+
}),
358+
assertSucceeded: vi.fn(() => {
359+
calls.push('receipt')
360+
}),
361+
}
362+
mockCreateSlackStreamController.mockResolvedValue(controller)
363+
mockExecuteWorkflowCore.mockImplementationOnce(async (options) => {
364+
const result = {
365+
success: true,
366+
status: 'completed',
367+
output: { content: 'Complete answer' },
368+
logs: [],
369+
}
370+
expect(options.callbacks).toBe(controller.callbacks)
371+
await options.finalizeDelivery(result)
372+
calls.push('workflow-result')
373+
return result
374+
})
375+
await executeWebhookJob({ ...legacyPayload, provider: 'slack' })
376+
expect(mockExecutionSnapshot.mock.calls[0]?.[0]).toMatchObject({
377+
agentEvents: true,
378+
includeThinking: false,
379+
includeToolCalls,
380+
})
381+
expect(mockExecutionSnapshot.mock.calls[0]?.[4]).toEqual(['mship_content'])
382+
expect(controller.finalize).toHaveBeenCalledTimes(1)
383+
expect(calls).toEqual(['delivery', 'receipt', 'workflow-result'])
384+
}
385+
)
386+
327387
it('restores a legacy queued webhook as its canonical system principal', async () => {
328388
mockExecuteWorkflowCore.mockResolvedValueOnce({
329389
success: true,

‎apps/sim/background/webhook-execution.ts‎

Lines changed: 14 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1120,6 +1120,12 @@ async function executeWebhookJobInternal(
11201120
})
11211121
: null
11221122

1123+
if (slackStreamConfig) {
1124+
metadata.agentEvents = true
1125+
metadata.includeThinking = slackStreamConfig.includeThinking
1126+
metadata.includeToolCalls = slackStreamConfig.includeToolCalls
1127+
}
1128+
11231129
const snapshot = new ExecutionSnapshot(
11241130
metadata,
11251131
workflowRecord,
@@ -1134,6 +1140,14 @@ async function executeWebhookJobInternal(
11341140
executionResult = await executeWorkflowCore({
11351141
snapshot,
11361142
callbacks: slackStreamController?.callbacks ?? {},
1143+
...(slackStreamController
1144+
? {
1145+
finalizeDelivery: async (result: ExecutionResult) => {
1146+
await slackStreamController.finalize(result)
1147+
slackStreamController.assertSucceeded()
1148+
},
1149+
}
1150+
: {}),
11371151
loggingSession,
11381152
trustedInitialResolvedSecretTraceProvenance:
11391153
resolvedSecretTraceRegistry.exportProvenanceForValue(triggerInput),
@@ -1151,11 +1165,6 @@ async function executeWebhookJobInternal(
11511165
}
11521166
throw error
11531167
}
1154-
if (slackStreamController) {
1155-
await slackStreamController.finalize(executionResult)
1156-
slackStreamController.assertSucceeded()
1157-
}
1158-
11591168
await handleExecutionResult(executionResult, {
11601169
loggingSession,
11611170
timeoutController,

‎apps/sim/executor/handlers/mothership/mothership-handler.test.ts‎

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1947,14 +1947,22 @@ describe('MothershipBlockHandler', () => {
19471947
context.metadata.agentEvents = agentEvents
19481948
fetchMock.mockResolvedValue(
19491949
createNdjsonResponse([
1950-
{ type: 'chunk', content: 'I will check.', turn: 'pending' },
1950+
{
1951+
type: 'agent_event',
1952+
v: 1,
1953+
event: { type: 'text_delta', text: 'I will check.', turn: 'pending' },
1954+
},
19511955
{ type: 'agent_event', event: { type: 'turn_end', turn: 'intermediate' } },
19521956
{ type: 'agent_event', event: { type: 'tool_call_start', id: 'lookup', name: 'Lookup' } },
19531957
{
19541958
type: 'agent_event',
19551959
event: { type: 'tool_call_end', id: 'lookup', name: 'Lookup', status: 'success' },
19561960
},
1957-
{ type: 'chunk', content: 'Final answer.', turn: 'pending' },
1961+
{
1962+
type: 'agent_event',
1963+
v: 1,
1964+
event: { type: 'text_delta', text: 'Final answer.', turn: 'pending' },
1965+
},
19581966
{ type: 'agent_event', event: { type: 'turn_end', turn: 'final' } },
19591967
{ type: 'final', data: { content: 'Final answer.' } },
19601968
])

‎apps/sim/executor/handlers/mothership/mothership-handler.ts‎

Lines changed: 12 additions & 26 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,10 @@
11
import { createLogger } from '@sim/logger'
22
import { generateId } from '@sim/utils/id'
33
import { isPlainRecord } from '@sim/utils/object'
4+
import type {
5+
MothershipExecuteResult,
6+
MothershipExecuteStreamEvent,
7+
} from '@/lib/api/contracts/mothership-chats'
48
import { generateInternalDelegationToken } from '@/lib/auth/internal'
59
import {
610
BILLING_ATTRIBUTION_HEADER,
@@ -53,11 +57,7 @@ import type {
5357
ResolvedSecretInputPath,
5458
ResolvedSecretTraceRegistry,
5559
} from '@/executor/utils/resolved-secret-trace-registry'
56-
import {
57-
type AgentStreamEvent,
58-
isAgentStreamEvent,
59-
type TextDeltaClassification,
60-
} from '@/providers/stream-events'
60+
import { type AgentStreamEvent, isAgentStreamEvent } from '@/providers/stream-events'
6161
import type { SerializedBlock } from '@/serializer/types'
6262

6363
const logger = createLogger('MothershipBlockHandler')
@@ -101,24 +101,6 @@ interface IndexedMothershipSkillContext {
101101
hasExplicitLabel: boolean
102102
}
103103

104-
type MothershipExecuteResult = {
105-
content?: string
106-
model?: string
107-
conversationId?: string
108-
tokens?: Record<string, unknown>
109-
toolCalls?: Array<Record<string, unknown>>
110-
cost?: unknown
111-
} & Partial<Record<typeof RESOLVED_SECRET_PROVENANCE_FIELD, unknown>>
112-
113-
type MothershipExecuteStreamEvent =
114-
| { type: 'heartbeat'; timestamp?: string }
115-
| { type: 'chunk'; content?: string; turn?: TextDeltaClassification }
116-
| { type: 'agent_event'; event: AgentStreamEvent }
117-
| { type: 'final'; data: MothershipExecuteResult }
118-
| ({ type: 'error'; error?: string } & Partial<
119-
Record<typeof RESOLVED_SECRET_PROVENANCE_FIELD, unknown>
120-
>)
121-
122104
function selectIndexedMothershipMcpTools(tools: unknown): IndexedMothershipMcpToolSelection[] {
123105
if (!Array.isArray(tools)) return []
124106

@@ -496,7 +478,7 @@ function formatMothershipBlockOutput(
496478
result: MothershipExecuteResult,
497479
conversationId: string
498480
): NormalizedBlockOutput {
499-
const formattedList = (result.toolCalls || []).map((tc: Record<string, unknown>) => ({
481+
const formattedList = (result.toolCalls || []).map((tc) => ({
500482
name: typeof tc.name === 'string' ? tc.name : String(tc.name ?? ''),
501483
...(typeof tc.status === 'string' ? { status: tc.status } : {}),
502484
arguments: (tc.arguments || tc.params || tc.input || {}) as Record<string, unknown>,
@@ -676,11 +658,15 @@ function createMothershipStreamingExecution(
676658
}
677659

678660
if (event.type === 'agent_event') {
679-
if (!isAgentStreamEvent(event.event)) {
661+
if ((event.v !== undefined && event.v !== 1) || !isAgentStreamEvent(event.event)) {
680662
throw new Error('Sim execution stream returned an invalid agent event')
681663
}
682664
if (options.agentEvents) controller.enqueue(event.event)
683-
else if (event.event.type === 'turn_end') {
665+
else if (event.event.type === 'text_delta') {
666+
if (event.event.turn === 'pending') pendingText += event.event.text
667+
else if (event.event.turn !== 'intermediate')
668+
controller.enqueue(encoder.encode(event.event.text))
669+
} else if (event.event.type === 'turn_end') {
684670
if (event.event.turn === 'final' && pendingText)
685671
controller.enqueue(encoder.encode(pendingText))
686672
pendingText = ''

0 commit comments

Comments
 (0)