Skip to content

Commit 5884a60

Browse files
committed
fix(slack): bound each long-answer append request
1 parent 66186b2 commit 5884a60

2 files changed

Lines changed: 78 additions & 11 deletions

File tree

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

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,7 @@ vi.mock('@/lib/webhooks/slack-stream-sessions', () => ({
4040
}))
4141

4242
import { ExecuteEventProjection } from '@/lib/mothership/request/lifecycle/execute-events'
43+
import type { SlackStreamChunk } from '@/lib/webhooks/slack-agent-api'
4344
import { SlackExecutionStreamController } from '@/lib/webhooks/slack-execution-stream'
4445
import type { SlackStreamResponseConfig } from '@/lib/webhooks/slack-stream-config'
4546
import { type AgentStreamEvent, createAgentEventReadableStream } from '@/providers/stream-events'
@@ -598,6 +599,65 @@ describe('SlackExecutionStreamController', () => {
598599
controller.assertSucceeded()
599600
})
600601

602+
it.each(['live', 'settled'] as const)(
603+
'delivers a long %s answer with each Slack append inside the request limit',
604+
async (delivery) => {
605+
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+
)
616+
const { controller } = await deliver(
617+
delivery === 'live'
618+
? [
619+
{ type: 'text_delta', text: answer, turn: 'pending' },
620+
{ type: 'turn_end', turn: 'final' },
621+
]
622+
: [],
623+
answer
624+
)
625+
controller.assertSucceeded()
626+
expect(sentText()).toBe(answer)
627+
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(3)
628+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(1)
629+
expect(mockStopSlackAgentStream).toHaveBeenCalledTimes(1)
630+
}
631+
)
632+
633+
it('does not replay acknowledged long-answer chunks after a later append fails', async () => {
634+
mockAppendSlackAgentStream.mockResolvedValueOnce(undefined)
635+
mockAppendSlackAgentStream.mockRejectedValueOnce(new Error('append acknowledgment lost'))
636+
const answer = `${'x'.repeat(25_000)}The complete ending.`
637+
const { controller } = await deliver(
638+
[
639+
{ type: 'text_delta', text: answer, turn: 'pending' },
640+
{ type: 'turn_end', turn: 'final' },
641+
],
642+
answer
643+
)
644+
expect(() => controller.assertSucceeded()).toThrow('append acknowledgment lost')
645+
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(2)
646+
expect(mockAppendSlackAgentStream.mock.calls.map((call) => call[3])).toEqual([
647+
[{ type: 'markdown_text', text: answer.slice(0, 12_000) }],
648+
[{ type: 'markdown_text', text: answer.slice(12_000, 24_000) }],
649+
])
650+
expect(mockStopSlackAgentStream).toHaveBeenCalledWith(
651+
expect.anything(),
652+
expect.anything(),
653+
expect.anything(),
654+
'suspended',
655+
undefined,
656+
undefined,
657+
[expect.objectContaining({ status: 'error' })]
658+
)
659+
})
660+
601661
it('reconciles a settled suffix without duplicating acknowledged text', async () => {
602662
const { controller } = await deliver(
603663
[

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

Lines changed: 18 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -87,13 +87,14 @@ export function resolveSlackReplyTarget(triggerInput: Record<string, unknown>):
8787
}
8888
}
8989

90-
function splitMarkdown(text: string): SlackStreamChunk[] {
91-
const chunks: SlackStreamChunk[] = []
92-
for (let offset = 0; offset < text.length; offset += SLACK_MARKDOWN_LIMIT) {
93-
chunks.push({
94-
type: 'markdown_text',
95-
text: text.slice(offset, offset + SLACK_MARKDOWN_LIMIT),
96-
})
90+
function splitMarkdown(text: string): string[] {
91+
const chunks: string[] = []
92+
for (let offset = 0; offset < text.length; ) {
93+
let end = Math.min(offset + SLACK_MARKDOWN_LIMIT, text.length)
94+
const lastCodeUnit = text.charCodeAt(end - 1)
95+
if (end < text.length && lastCodeUnit >= 0xd800 && lastCodeUnit <= 0xdbff) end--
96+
chunks.push(text.slice(offset, end))
97+
offset = end
9798
}
9899
return chunks
99100
}
@@ -156,6 +157,14 @@ class SlackInvocationStream {
156157
await appendSlackAgentStream(this.token, this.channel!, this.ts!, chunks, this.signal)
157158
}
158159

160+
/** Slack limits the whole append request, not each chunk within its array. */
161+
private async appendAnswerText(text: string): Promise<void> {
162+
for (const chunk of splitMarkdown(text)) {
163+
await this.append([{ type: 'markdown_text', text: chunk }])
164+
this.acknowledgedAnswer += chunk
165+
}
166+
}
167+
159168
private async flushAnswer(force: boolean): Promise<void> {
160169
if (
161170
this.failure ||
@@ -165,8 +174,7 @@ class SlackInvocationStream {
165174
return
166175
const projected = await this.projectLiveText(this.answerBuffer)
167176
if (projected === null) return
168-
if (projected) await this.append(splitMarkdown(projected))
169-
this.acknowledgedAnswer += projected
177+
await this.appendAnswerText(projected)
170178
this.answerBuffer = ''
171179
}
172180

@@ -267,8 +275,7 @@ class SlackInvocationStream {
267275
throw new Error('Slack delivery failed: settled output differs from acknowledged answer')
268276
}
269277
const remaining = text.slice(this.acknowledgedAnswer.length)
270-
if (remaining) await this.append(splitMarkdown(remaining))
271-
this.acknowledgedAnswer = text
278+
await this.appendAnswerText(remaining)
272279
this.answerBuffer = ''
273280
})
274281
}

0 commit comments

Comments
 (0)