Skip to content

Commit 6cd8f16

Browse files
committed
fix(slack): continue long streams across message limits
1 parent 5884a60 commit 6cd8f16

2 files changed

Lines changed: 262 additions & 70 deletions

File tree

‎apps/sim/lib/webhooks/slack-execution-stream.test.ts‎

Lines changed: 163 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -475,6 +475,51 @@ describe('SlackExecutionStreamController', () => {
475475
.join('')
476476
}
477477

478+
function enforceSlackMessageLimit() {
479+
const messages = new Map<
480+
string,
481+
{ text: string; chunks: SlackStreamChunk[]; stopped: boolean }
482+
>()
483+
mockStartSlackAgentStream.mockImplementation(async (_token, _target, chunks) => {
484+
const ts = `message-${messages.size}`
485+
messages.set(ts, { text: '', chunks: [...chunks], stopped: false })
486+
return { channel: 'C123', ts }
487+
})
488+
mockAppendSlackAgentStream.mockImplementation(
489+
async (_token, _channel, ts: string, chunks: SlackStreamChunk[]) => {
490+
const message = messages.get(ts)!
491+
if (message.stopped) throw new Error('message_not_in_streaming_state')
492+
const text = chunks
493+
.filter((chunk) => chunk.type === 'markdown_text')
494+
.map((chunk) => chunk.text)
495+
.join('')
496+
if (message.text.length + text.length > 12_000) {
497+
throw new Error('Slack chat.appendStream: msg_too_long')
498+
}
499+
expect(text.isWellFormed()).toBe(true)
500+
message.text += text
501+
message.chunks.push(...chunks)
502+
}
503+
)
504+
mockStopSlackAgentStream.mockImplementation(
505+
async (
506+
_token,
507+
_channel,
508+
ts: string,
509+
_status,
510+
_signal,
511+
_blocks,
512+
chunks: SlackStreamChunk[]
513+
) => {
514+
const message = messages.get(ts)!
515+
expect(message.stopped).toBe(false)
516+
message.stopped = true
517+
message.chunks.push(...chunks)
518+
}
519+
)
520+
return messages
521+
}
522+
478523
it.each([true, false])(
479524
'delivers identical Agent and Mship thinking updates with thinking enabled: %s',
480525
async (includeThinking) => {
@@ -600,19 +645,10 @@ describe('SlackExecutionStreamController', () => {
600645
})
601646

602647
it.each(['live', 'settled'] as const)(
603-
'delivers a long %s answer with each Slack append inside the request limit',
648+
'delivers a long %s answer across messages within Slack’s cumulative text limit',
604649
async (delivery) => {
605650
const answer = `${'x'.repeat(11_999)}🚀${'y'.repeat(17_000)}The complete ending.`
606-
mockAppendSlackAgentStream.mockImplementation(
607-
async (_token, _channel, _ts, chunks: SlackStreamChunk[]) => {
608-
const text = chunks
609-
.filter((chunk) => chunk.type === 'markdown_text')
610-
.map((chunk) => chunk.text)
611-
.join('')
612-
if (text.length > 12_000) throw new Error('Slack chat.appendStream: msg_too_long')
613-
expect(text.isWellFormed()).toBe(true)
614-
}
615-
)
651+
const messages = enforceSlackMessageLimit()
616652
const { controller } = await deliver(
617653
delivery === 'live'
618654
? [
@@ -624,12 +660,126 @@ describe('SlackExecutionStreamController', () => {
624660
)
625661
controller.assertSucceeded()
626662
expect(sentText()).toBe(answer)
663+
expect([...messages.values()].map((message) => message.text).join('')).toBe(answer)
664+
expect([...messages.values()].every((message) => message.stopped)).toBe(true)
627665
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(3)
628-
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(1)
629-
expect(mockStopSlackAgentStream).toHaveBeenCalledTimes(1)
666+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(3)
667+
expect(mockStopSlackAgentStream).toHaveBeenCalledTimes(3)
668+
expect(mockStopSlackAgentStream.mock.calls.map((call) => call[2])).toEqual([
669+
...messages.keys(),
670+
])
671+
for (const call of mockStartSlackAgentStream.mock.calls) {
672+
expect(call[1]).toMatchObject({ channel: 'C123', threadTs: '1700000000.000001' })
673+
}
630674
}
631675
)
632676

677+
it('counts commentary across turns and pairs tools on their original message after rollover', async () => {
678+
const messages = enforceSlackMessageLimit()
679+
const commentary = 'c'.repeat(11_800)
680+
const answer = `${'a'.repeat(15_000)}The complete ending.`
681+
const events: AgentStreamEvent[] = [
682+
{ type: 'text_delta', text: commentary, turn: 'pending' },
683+
{ type: 'turn_end', turn: 'intermediate' },
684+
{ type: 'tool_call_start', id: 'a', name: 'read' },
685+
{ type: 'tool_call_start', id: 'b', name: 'search' },
686+
...Array.from(
687+
{ length: 100 },
688+
(_, index): AgentStreamEvent => ({
689+
type: 'text_delta',
690+
text: answer.slice(index * 150, (index + 1) * 150),
691+
turn: 'pending',
692+
})
693+
),
694+
{ type: 'tool_call_end', id: 'b', name: 'search', status: 'success' },
695+
{ type: 'tool_call_end', id: 'a', name: 'read', status: 'success' },
696+
{ type: 'tool_call_start', id: 'c', name: 'read' },
697+
{ type: 'tool_call_end', id: 'c', name: 'read', status: 'success' },
698+
{ type: 'turn_end', turn: 'final' },
699+
]
700+
const { controller } = await deliver(events, answer)
701+
controller.assertSucceeded()
702+
expect([...messages.values()].map((message) => message.text).join('')).toBe(commentary + answer)
703+
expect(messages.size).toBe(3)
704+
for (const message of messages.values()) {
705+
expect(message.stopped).toBe(true)
706+
const tasks = message.chunks.filter((chunk) => chunk.type === 'task_update')
707+
for (const taskId of new Set(tasks.map((task) => task.id))) {
708+
expect(tasks.filter((task) => task.id === taskId).map((task) => task.status)).toEqual([
709+
'in_progress',
710+
'complete',
711+
])
712+
}
713+
}
714+
const first = messages.get('message-0')!
715+
expect(
716+
first.chunks.filter((chunk) => chunk.type === 'task_update' && chunk.id.includes('-tool-'))
717+
).toHaveLength(4)
718+
expect(mockSetSlackAgentSessionStatus.mock.calls.at(-1)?.[2]).toBe('active')
719+
})
720+
721+
it('cleans up every continuation on cancellation and marks tools on the right message', async () => {
722+
const messages = enforceSlackMessageLimit()
723+
const { controller } = await createController()
724+
const pump = createAgentStreamPump({
725+
source: createAgentEventReadableStream([
726+
{ type: 'tool_call_start', id: 'running', name: 'read' },
727+
{ type: 'text_delta', text: 'a'.repeat(25_000), turn: 'pending' },
728+
]),
729+
streamFormat: 'agent-events-v1',
730+
})
731+
await Promise.all([
732+
controller.callbacks.onStream?.({
733+
blockId: 'agent',
734+
executionOrder: 1,
735+
stream: pump.textStream!,
736+
subscribe: pump.subscribe,
737+
}),
738+
pump.run(),
739+
])
740+
await controller.finalize({ success: false, status: 'cancelled', output: {} })
741+
expect(messages.size).toBe(3)
742+
expect([...messages.values()].every((message) => message.stopped)).toBe(true)
743+
const toolStops = mockStopSlackAgentStream.mock.calls.filter((call) =>
744+
call[6].some(
745+
(chunk: SlackStreamChunk) =>
746+
chunk.type === 'task_update' && chunk.id.endsWith('-tool-running')
747+
)
748+
)
749+
expect(toolStops).toHaveLength(1)
750+
expect(toolStops[0][2]).toBe('message-0')
751+
expect(toolStops[0][6][0].status).toBe('error')
752+
expect(mockUnregisterSlackStreamSession).toHaveBeenCalledTimes(1)
753+
})
754+
755+
it('still stops the remaining messages when one continuation stop fails', async () => {
756+
const messages = enforceSlackMessageLimit()
757+
mockStopSlackAgentStream.mockRejectedValueOnce(new Error('stop outcome uncertain'))
758+
const { controller } = await deliver([], 'a'.repeat(25_000))
759+
expect(() => controller.assertSucceeded()).toThrow('stop outcome uncertain')
760+
expect(mockStopSlackAgentStream.mock.calls.map((call) => call[2])).toEqual([...messages.keys()])
761+
expect(messages.get('message-1')!.stopped).toBe(true)
762+
expect(messages.get('message-2')!.stopped).toBe(true)
763+
expect(mockSetSlackAgentSessionStatus.mock.calls.at(-1)?.[2]).toBe('suspended')
764+
await controller.finalize({ success: true, output: {} })
765+
expect(mockStopSlackAgentStream).toHaveBeenCalledTimes(3)
766+
})
767+
768+
it('does not retry an uncertain continuation start or duplicate the first part', async () => {
769+
const messages = enforceSlackMessageLimit()
770+
mockStartSlackAgentStream
771+
.mockImplementationOnce(async (_token, _target, chunks) => {
772+
messages.set('message-0', { text: '', chunks, stopped: false })
773+
return { channel: 'C123', ts: 'message-0' }
774+
})
775+
.mockRejectedValueOnce(new Error('continuation start outcome uncertain'))
776+
const { controller } = await deliver([], 'a'.repeat(25_000))
777+
expect(() => controller.assertSucceeded()).toThrow('continuation start outcome uncertain')
778+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(2)
779+
expect(sentText()).toBe('a'.repeat(12_000))
780+
expect(messages.get('message-0')!.stopped).toBe(true)
781+
})
782+
633783
it('does not replay acknowledged long-answer chunks after a later append fails', async () => {
634784
mockAppendSlackAgentStream.mockResolvedValueOnce(undefined)
635785
mockAppendSlackAgentStream.mockRejectedValueOnce(new Error('append acknowledgment lost'))

0 commit comments

Comments
 (0)