Skip to content

Commit a09adba

Browse files
committed
fix(slack): clean up long-response continuations
1 parent d0ddd33 commit a09adba

2 files changed

Lines changed: 217 additions & 52 deletions

File tree

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

Lines changed: 127 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -309,7 +309,7 @@ describe('SlackExecutionStreamController', () => {
309309
'xoxb-token',
310310
'C123',
311311
'1700000001.000002',
312-
[{ type: 'markdown_text', text: 'Once upon a time' }],
312+
[{ type: 'markdown_text', text: 'Once upon a ' }],
313313
undefined
314314
)
315315
})
@@ -367,7 +367,8 @@ describe('SlackExecutionStreamController', () => {
367367
title: 'Gmail Send Email',
368368
status: 'complete',
369369
}),
370-
{ type: 'markdown_text', text: 'Selected answer' },
370+
{ type: 'markdown_text', text: 'Selected ' },
371+
{ type: 'markdown_text', text: 'answer' },
371372
])
372373
)
373374
expect(appendedChunks).not.toContainEqual({
@@ -469,8 +470,18 @@ describe('SlackExecutionStreamController', () => {
469470
}
470471

471472
function sentText(): string {
472-
return mockAppendSlackAgentStream.mock.calls
473-
.flatMap((call) => call[3])
473+
return [
474+
...mockStartSlackAgentStream.mock.calls.map((call, index) => ({
475+
order: mockStartSlackAgentStream.mock.invocationCallOrder[index],
476+
chunks: call[2],
477+
})),
478+
...mockAppendSlackAgentStream.mock.calls.map((call, index) => ({
479+
order: mockAppendSlackAgentStream.mock.invocationCallOrder[index],
480+
chunks: call[3],
481+
})),
482+
]
483+
.sort((a, b) => a.order - b.order)
484+
.flatMap((call) => call.chunks)
474485
.filter((chunk) => chunk.type === 'markdown_text')
475486
.map((chunk) => chunk.text)
476487
.join('')
@@ -482,8 +493,16 @@ describe('SlackExecutionStreamController', () => {
482493
{ text: string; chunks: SlackStreamChunk[]; stopped: boolean }
483494
>()
484495
mockStartSlackAgentStream.mockImplementation(async (_token, _target, chunks) => {
496+
const text = chunks
497+
.filter((chunk: SlackStreamChunk) => chunk.type === 'markdown_text')
498+
.map((chunk: SlackStreamChunk) => (chunk.type === 'markdown_text' ? chunk.text : ''))
499+
.join('')
500+
if (text.length > limit) {
501+
throw new SlackDeliveryError('chat.startStream', 'rejected', 'msg_too_long', 200)
502+
}
503+
expect(text.isWellFormed()).toBe(true)
485504
const ts = `message-${messages.size}`
486-
messages.set(ts, { text: '', chunks: [...chunks], stopped: false })
505+
messages.set(ts, { text, chunks: [...chunks], stopped: false })
487506
return { channel: 'C123', ts }
488507
})
489508
mockAppendSlackAgentStream.mockImplementation(
@@ -663,7 +682,6 @@ describe('SlackExecutionStreamController', () => {
663682
expect(sentText()).toBe(answer)
664683
expect([...messages.values()].map((message) => message.text).join('')).toBe(answer)
665684
expect([...messages.values()].every((message) => message.stopped)).toBe(true)
666-
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(3)
667685
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(3)
668686
expect(mockStopSlackAgentStream).toHaveBeenCalledTimes(3)
669687
expect(mockStopSlackAgentStream.mock.calls.map((call) => call[2])).toEqual([
@@ -677,8 +695,8 @@ describe('SlackExecutionStreamController', () => {
677695

678696
it('continues a definitively rejected 11938 + 62 boundary append without replaying accepted text', async () => {
679697
const messages = enforceSlackMessageLimit(11_999)
680-
const prefix = 'a'.repeat(11_938)
681-
const suffix = `${'b'.repeat(3_000)}Complete ending.`
698+
const prefix = `${'a'.repeat(11_937)} `
699+
const suffix = `${'b'.repeat(61)} ${'b'.repeat(2_938)}Complete ending.`
682700
const { controller } = await deliver(
683701
[
684702
{ type: 'tool_call_start', id: 'read', name: 'read' },
@@ -707,6 +725,69 @@ describe('SlackExecutionStreamController', () => {
707725
expect(toolChunks.map((chunk) => chunk.status)).toEqual(['in_progress', 'complete'])
708726
})
709727

728+
it('keeps a word split across live deltas together when continuing to another reply', async () => {
729+
const messages = enforceSlackMessageLimit()
730+
const prefix = `${'Calm seas. '.repeat(1_090)}Mara `
731+
const { controller } = await deliver(
732+
[
733+
{ type: 'tool_call_start', id: 'read', name: 'read' },
734+
{ type: 'text_delta', text: `${prefix}lau`, turn: 'pending' },
735+
{ type: 'text_delta', text: 'gh', turn: 'pending' },
736+
{ type: 'text_delta', text: 'ed once, sharply, with relief.', turn: 'pending' },
737+
{ type: 'tool_call_end', id: 'read', name: 'read', status: 'success' },
738+
{ type: 'turn_end', turn: 'final' },
739+
],
740+
`${prefix}laughed once, sharply, with relief.`
741+
)
742+
controller.assertSucceeded()
743+
const parts = [...messages.values()]
744+
expect(parts.map((part) => part.text)).toEqual([prefix, 'laughed once, sharply, with relief.'])
745+
expect(parts[1].chunks.every((chunk) => chunk.type === 'markdown_text')).toBe(true)
746+
expect(parts.every((part) => part.stopped)).toBe(true)
747+
expect(parts[0].chunks.filter((chunk) => chunk.type === 'task_update')).toEqual([
748+
expect.objectContaining({ title: 'Running', status: 'in_progress' }),
749+
expect.objectContaining({ id: 'sim-execution-1-1-tool-read', status: 'in_progress' }),
750+
expect.objectContaining({ id: 'sim-execution-1-1-tool-read', status: 'complete' }),
751+
expect.objectContaining({ title: 'Running', status: 'complete' }),
752+
])
753+
})
754+
755+
it.each(['\n\n', '. ', ' '])(
756+
'breaks settled prose at the available %j boundary',
757+
async (separator) => {
758+
const messages = enforceSlackMessageLimit()
759+
const first = `${'x'.repeat(11_500)}${separator}`
760+
const second = `${'word '.repeat(250)}Complete ending.`
761+
const { controller } = await deliver([], first + second)
762+
controller.assertSucceeded()
763+
const parts = [...messages.values()]
764+
expect(parts.map((part) => part.text).join('')).toBe(first + second)
765+
if (separator !== ' ') expect(parts[0].text).toBe(first)
766+
expect(parts[0].text).toMatch(/\s$/)
767+
expect(parts[1].text).toMatch(/^word /)
768+
expect(parts[1].chunks.every((chunk) => chunk.type === 'markdown_text')).toBe(true)
769+
}
770+
)
771+
772+
it('preserves all prose and clean word boundaries through many incremental deltas', async () => {
773+
const messages = enforceSlackMessageLimit()
774+
const answer = `${'Mara laughed, then watched the sea.\n\n'.repeat(900)}The end.`
775+
const events: AgentStreamEvent[] = []
776+
for (let offset = 0; offset < answer.length; offset += 17) {
777+
events.push({ type: 'text_delta', text: answer.slice(offset, offset + 17), turn: 'pending' })
778+
}
779+
events.push({ type: 'turn_end', turn: 'final' })
780+
const { controller } = await deliver(events, answer)
781+
controller.assertSucceeded()
782+
const parts = [...messages.values()]
783+
expect(parts.length).toBeGreaterThan(2)
784+
expect(parts.map((part) => part.text).join('')).toBe(answer)
785+
for (const part of parts.slice(0, -1)) expect(part.text).toMatch(/\s$/)
786+
for (const part of parts.slice(1)) {
787+
expect(part.chunks.every((chunk) => chunk.type === 'markdown_text')).toBe(true)
788+
}
789+
})
790+
710791
it('reduces a definitively oversized first append and still delivers all final-only text', async () => {
711792
const messages = enforceSlackMessageLimit(11_999)
712793
const answer = '🚀'.repeat(13_000)
@@ -743,28 +824,31 @@ describe('SlackExecutionStreamController', () => {
743824
expect(messages.get('message-0')!.stopped).toBe(true)
744825
})
745826

746-
it('stops recovery if the continuation append has an uncertain outcome', async () => {
827+
it('stops recovery if the continuation start has an uncertain outcome', async () => {
747828
const messages = enforceSlackMessageLimit()
748-
const prefix = 'a'.repeat(11_938)
829+
const prefix = `${'a'.repeat(11_937)} `
830+
mockStartSlackAgentStream
831+
.mockImplementationOnce(mockStartSlackAgentStream.getMockImplementation()!)
832+
.mockRejectedValueOnce(
833+
new SlackDeliveryError('chat.startStream', 'uncertain', 'transport_failed')
834+
)
749835
mockAppendSlackAgentStream
750836
.mockImplementationOnce(async () => {
751837
messages.get('message-0')!.text = prefix
752838
})
753839
.mockRejectedValueOnce(
754840
new SlackDeliveryError('chat.appendStream', 'rejected', 'msg_too_long')
755841
)
756-
.mockRejectedValueOnce(
757-
new SlackDeliveryError('chat.appendStream', 'uncertain', 'transport_failed')
758-
)
842+
const suffix = `${'b'.repeat(61)} ${'b'.repeat(88)}`
759843
const { controller } = await deliver(
760844
[
761845
{ type: 'text_delta', text: prefix, turn: 'pending' },
762-
{ type: 'text_delta', text: 'b'.repeat(150), turn: 'pending' },
846+
{ type: 'text_delta', text: suffix, turn: 'pending' },
763847
],
764-
prefix + 'b'.repeat(150)
848+
prefix + suffix
765849
)
766850
expect(() => controller.assertSucceeded()).toThrow('transport_failed')
767-
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(3)
851+
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(2)
768852
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(2)
769853
expect([...messages.values()].every((message) => message.stopped)).toBe(true)
770854
})
@@ -871,12 +955,15 @@ describe('SlackExecutionStreamController', () => {
871955
const { controller } = await deliver([], 'a'.repeat(25_000))
872956
expect(() => controller.assertSucceeded()).toThrow('continuation start outcome uncertain')
873957
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(2)
874-
expect(sentText()).toBe('a'.repeat(12_000))
958+
expect(messages.get('message-0')!.text).toBe('a'.repeat(12_000))
875959
expect(messages.get('message-0')!.stopped).toBe(true)
876960
})
877961

878962
it('does not replay acknowledged long-answer chunks after a later append fails', async () => {
879-
mockAppendSlackAgentStream.mockResolvedValueOnce(undefined)
963+
const messages = enforceSlackMessageLimit()
964+
mockAppendSlackAgentStream.mockImplementationOnce(
965+
mockAppendSlackAgentStream.getMockImplementation()!
966+
)
880967
mockAppendSlackAgentStream.mockRejectedValueOnce(new Error('append acknowledgment lost'))
881968
const answer = `${'x'.repeat(25_000)}The complete ending.`
882969
const { controller } = await deliver(
@@ -890,8 +977,12 @@ describe('SlackExecutionStreamController', () => {
890977
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(2)
891978
expect(mockAppendSlackAgentStream.mock.calls.map((call) => call[3])).toEqual([
892979
[{ type: 'markdown_text', text: answer.slice(0, 12_000) }],
893-
[{ type: 'markdown_text', text: answer.slice(12_000, 24_000) }],
980+
[{ type: 'markdown_text', text: 'ending.' }],
894981
])
982+
expect([...messages.values()].map((message) => message.text).join('')).toBe(
983+
answer.slice(0, -'ending.'.length)
984+
)
985+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(3)
895986
expect(mockStopSlackAgentStream).toHaveBeenCalledWith(
896987
expect.anything(),
897988
expect.anything(),
@@ -903,6 +994,23 @@ describe('SlackExecutionStreamController', () => {
903994
)
904995
})
905996

997+
it('bounds continuation recovery when even a single character is explicitly rejected', async () => {
998+
const messages = enforceSlackMessageLimit()
999+
mockStartSlackAgentStream
1000+
.mockImplementationOnce(mockStartSlackAgentStream.getMockImplementation()!)
1001+
.mockRejectedValue(new SlackDeliveryError('chat.startStream', 'rejected', 'msg_too_long'))
1002+
const prefix = 'word '.repeat(2_400)
1003+
const { controller } = await deliver([], `${prefix}abcd`)
1004+
expect(() => controller.assertSucceeded()).toThrow('msg_too_long')
1005+
expect(mockStartSlackAgentStream.mock.calls.slice(1).map((call) => call[2])).toEqual([
1006+
[{ type: 'markdown_text', text: 'abcd' }],
1007+
[{ type: 'markdown_text', text: 'ab' }],
1008+
[{ type: 'markdown_text', text: 'a' }],
1009+
])
1010+
expect([...messages.values()].map((message) => message.text)).toEqual([prefix])
1011+
expect(messages.get('message-0')!.stopped).toBe(true)
1012+
})
1013+
9061014
it('reconciles a settled suffix without duplicating acknowledged text', async () => {
9071015
const { controller } = await deliver(
9081016
[

0 commit comments

Comments
 (0)