Skip to content

Commit f23589f

Browse files
committed
fix(slack): recover explicitly rejected size overflows
1 parent 6cd8f16 commit f23589f

5 files changed

Lines changed: 162 additions & 20 deletions

File tree

‎apps/sim/lib/webhooks/slack-agent-api.test.ts‎

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ import {
88
startSlackAgentStream,
99
stopSlackAgentStream,
1010
} from '@/lib/webhooks/slack-agent-api'
11+
import { SlackDeliveryError } from '@/lib/webhooks/slack-delivery-error'
1112

1213
afterEach(() => {
1314
vi.unstubAllGlobals()
@@ -184,6 +185,26 @@ describe('Slack agent API transport', () => {
184185
).rejects.toThrow('missing_scope')
185186
})
186187

188+
it('exposes an explicit size rejection to the delivery controller without a transport retry', async () => {
189+
const fetchMock = vi
190+
.fn()
191+
.mockResolvedValue(
192+
new Response(JSON.stringify({ ok: false, error: 'msg_too_long' }), { status: 200 })
193+
)
194+
vi.stubGlobal('fetch', fetchMock)
195+
const appended = appendSlackAgentStream('token', 'C1', '1.2', [
196+
{ type: 'markdown_text', text: 'undelivered suffix' },
197+
])
198+
await expect(appended).rejects.toBeInstanceOf(SlackDeliveryError)
199+
await expect(appended).rejects.toMatchObject({
200+
method: 'chat.appendStream',
201+
outcome: 'rejected',
202+
code: 'msg_too_long',
203+
httpStatus: 200,
204+
})
205+
expect(fetchMock).toHaveBeenCalledTimes(1)
206+
})
207+
187208
it('fails fast when Slack does not recognize the stop-event subscription', async () => {
188209
vi.stubGlobal(
189210
'fetch',

‎apps/sim/lib/webhooks/slack-agent-api.ts‎

Lines changed: 1 addition & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
import { toError } from '@sim/utils/errors'
22
import { isRecordLike } from '@sim/utils/object'
33
import { readResponseJsonWithLimit } from '@/lib/core/utils/stream-limits'
4+
import { SlackDeliveryError } from '@/lib/webhooks/slack-delivery-error'
45

56
export type SlackStreamChunk =
67
| { type: 'markdown_text'; text: string }
@@ -31,19 +32,6 @@ interface SlackStreamTarget {
3132
recipientTeamId?: string
3233
}
3334

34-
/** A missing acknowledgment does not establish that Slack rejected a write. */
35-
export class SlackDeliveryError extends Error {
36-
constructor(
37-
readonly method: string,
38-
readonly outcome: 'rejected' | 'uncertain',
39-
readonly code: string,
40-
readonly httpStatus?: number
41-
) {
42-
super(`Slack ${method}: ${code} (delivery ${outcome})`)
43-
this.name = 'SlackDeliveryError'
44-
}
45-
}
46-
4735
/** Slack acknowledgments can include the full accumulated message, including task cards. */
4836
const MAX_SLACK_RESPONSE_BYTES = 4 * 1024 * 1024
4937

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,12 @@
1+
/** A missing acknowledgment does not establish that Slack rejected a write. */
2+
export class SlackDeliveryError extends Error {
3+
constructor(
4+
readonly method: string,
5+
readonly outcome: 'rejected' | 'uncertain',
6+
readonly code: string,
7+
readonly httpStatus?: number
8+
) {
9+
super(`Slack ${method}: ${code} (delivery ${outcome})`)
10+
this.name = 'SlackDeliveryError'
11+
}
12+
}

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

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

4242
import { ExecuteEventProjection } from '@/lib/mothership/request/lifecycle/execute-events'
4343
import type { SlackStreamChunk } from '@/lib/webhooks/slack-agent-api'
44+
import { SlackDeliveryError } from '@/lib/webhooks/slack-delivery-error'
4445
import { SlackExecutionStreamController } from '@/lib/webhooks/slack-execution-stream'
4546
import type { SlackStreamResponseConfig } from '@/lib/webhooks/slack-stream-config'
4647
import { type AgentStreamEvent, createAgentEventReadableStream } from '@/providers/stream-events'
@@ -475,7 +476,7 @@ describe('SlackExecutionStreamController', () => {
475476
.join('')
476477
}
477478

478-
function enforceSlackMessageLimit() {
479+
function enforceSlackMessageLimit(limit = 12_000) {
479480
const messages = new Map<
480481
string,
481482
{ text: string; chunks: SlackStreamChunk[]; stopped: boolean }
@@ -493,8 +494,8 @@ describe('SlackExecutionStreamController', () => {
493494
.filter((chunk) => chunk.type === 'markdown_text')
494495
.map((chunk) => chunk.text)
495496
.join('')
496-
if (message.text.length + text.length > 12_000) {
497-
throw new Error('Slack chat.appendStream: msg_too_long')
497+
if (message.text.length + text.length > limit) {
498+
throw new SlackDeliveryError('chat.appendStream', 'rejected', 'msg_too_long', 200)
498499
}
499500
expect(text.isWellFormed()).toBe(true)
500501
message.text += text
@@ -674,6 +675,100 @@ describe('SlackExecutionStreamController', () => {
674675
}
675676
)
676677

678+
it('continues a definitively rejected 11938 + 62 boundary append without replaying accepted text', async () => {
679+
const messages = enforceSlackMessageLimit(11_999)
680+
const prefix = 'a'.repeat(11_938)
681+
const suffix = `${'b'.repeat(3_000)}Complete ending.`
682+
const { controller } = await deliver(
683+
[
684+
{ type: 'tool_call_start', id: 'read', name: 'read' },
685+
{ type: 'text_delta', text: prefix, turn: 'pending' },
686+
{ type: 'text_delta', text: suffix, turn: 'pending' },
687+
{ type: 'tool_call_end', id: 'read', name: 'read', status: 'success' },
688+
{ type: 'turn_end', turn: 'final' },
689+
],
690+
prefix + suffix
691+
)
692+
controller.assertSucceeded()
693+
expect([...messages.values()].map((message) => message.text)).toEqual([prefix, suffix])
694+
const rejected = mockAppendSlackAgentStream.mock.calls.filter(
695+
(call) =>
696+
call[2] === 'message-0' &&
697+
call[3].some(
698+
(chunk: SlackStreamChunk) =>
699+
chunk.type === 'markdown_text' && chunk.text === suffix.slice(0, 62)
700+
)
701+
)
702+
expect(rejected).toHaveLength(1)
703+
expect([...messages.values()].every((message) => message.stopped)).toBe(true)
704+
const toolChunks = messages
705+
.get('message-0')!
706+
.chunks.filter((chunk) => chunk.type === 'task_update' && chunk.id.endsWith('-tool-read'))
707+
expect(toolChunks.map((chunk) => chunk.status)).toEqual(['in_progress', 'complete'])
708+
})
709+
710+
it('reduces a definitively oversized first append and still delivers all final-only text', async () => {
711+
const messages = enforceSlackMessageLimit(11_999)
712+
const answer = '🚀'.repeat(13_000)
713+
const { controller } = await deliver([], answer)
714+
controller.assertSucceeded()
715+
expect([...messages.values()].map((message) => message.text).join('')).toBe(answer)
716+
expect([...messages.values()].every((message) => message.text && message.stopped)).toBe(true)
717+
})
718+
719+
it.each([
720+
new SlackDeliveryError('chat.appendStream', 'uncertain', 'msg_too_long'),
721+
new SlackDeliveryError('chat.appendStream', 'rejected', 'ratelimited', 429),
722+
new SlackDeliveryError('chat.appendStream', 'rejected', 'missing_scope', 200),
723+
])('does not replay ambiguous or non-size append errors: %s', async (error) => {
724+
const messages = enforceSlackMessageLimit()
725+
mockAppendSlackAgentStream.mockRejectedValueOnce(error)
726+
const { controller } = await deliver([{ type: 'text_delta', text: 'answer' }], 'answer')
727+
expect(() => controller.assertSucceeded()).toThrow(error)
728+
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(1)
729+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(1)
730+
expect(messages.get('message-0')!.stopped).toBe(true)
731+
})
732+
733+
it('bounds size recovery when even a single character is rejected', async () => {
734+
const messages = enforceSlackMessageLimit(0)
735+
const { controller } = await deliver([], 'abcd')
736+
expect(() => controller.assertSucceeded()).toThrow('msg_too_long')
737+
expect(mockAppendSlackAgentStream.mock.calls.map((call) => call[3])).toEqual([
738+
[{ type: 'markdown_text', text: 'abcd' }],
739+
[{ type: 'markdown_text', text: 'ab' }],
740+
[{ type: 'markdown_text', text: 'a' }],
741+
])
742+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(1)
743+
expect(messages.get('message-0')!.stopped).toBe(true)
744+
})
745+
746+
it('stops recovery if the continuation append has an uncertain outcome', async () => {
747+
const messages = enforceSlackMessageLimit()
748+
const prefix = 'a'.repeat(11_938)
749+
mockAppendSlackAgentStream
750+
.mockImplementationOnce(async () => {
751+
messages.get('message-0')!.text = prefix
752+
})
753+
.mockRejectedValueOnce(
754+
new SlackDeliveryError('chat.appendStream', 'rejected', 'msg_too_long')
755+
)
756+
.mockRejectedValueOnce(
757+
new SlackDeliveryError('chat.appendStream', 'uncertain', 'transport_failed')
758+
)
759+
const { controller } = await deliver(
760+
[
761+
{ type: 'text_delta', text: prefix, turn: 'pending' },
762+
{ type: 'text_delta', text: 'b'.repeat(150), turn: 'pending' },
763+
],
764+
prefix + 'b'.repeat(150)
765+
)
766+
expect(() => controller.assertSucceeded()).toThrow('transport_failed')
767+
expect(mockAppendSlackAgentStream).toHaveBeenCalledTimes(3)
768+
expect(mockStartSlackAgentStream).toHaveBeenCalledTimes(2)
769+
expect([...messages.values()].every((message) => message.stopped)).toBe(true)
770+
})
771+
677772
it('counts commentary across turns and pairs tools on their original message after rollover', async () => {
678773
const messages = enforceSlackMessageLimit()
679774
const commentary = 'c'.repeat(11_800)

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

Lines changed: 30 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import {
1212
startSlackAgentStream,
1313
stopSlackAgentStream,
1414
} from '@/lib/webhooks/slack-agent-api'
15+
import { SlackDeliveryError } from '@/lib/webhooks/slack-delivery-error'
1516
import type {
1617
SlackStreamOutputConfig,
1718
SlackStreamResponseConfig,
@@ -189,10 +190,35 @@ class SlackInvocationStream {
189190
let end = Math.min(offset + SLACK_MESSAGE_TEXT_LIMIT - message.textLength, text.length)
190191
const lastCodeUnit = text.charCodeAt(end - 1)
191192
if (end < text.length && lastCodeUnit >= 0xd800 && lastCodeUnit <= 0xdbff) end--
192-
const chunk = text.slice(offset, end)
193-
await this.append([{ type: 'markdown_text', text: chunk }], message)
194-
this.acknowledgedAnswer += chunk
195-
offset = end
193+
for (;;) {
194+
const chunk = text.slice(offset, end)
195+
try {
196+
await this.append([{ type: 'markdown_text', text: chunk }], message)
197+
this.acknowledgedAnswer += chunk
198+
offset = end
199+
break
200+
} catch (error) {
201+
/** Only an explicit size rejection proves this text was not appended. */
202+
if (
203+
!(error instanceof SlackDeliveryError) ||
204+
error.method !== 'chat.appendStream' ||
205+
error.outcome !== 'rejected' ||
206+
error.code !== 'msg_too_long'
207+
) {
208+
throw error
209+
}
210+
if (message.textLength > 0) {
211+
message = await this.ensureStarted(true)
212+
} else {
213+
/** A rejected first append shrinks on each attempt, stopping at one code point. */
214+
let smallerEnd = offset + Math.floor((end - offset) / 2)
215+
const lastCodeUnit = text.charCodeAt(smallerEnd - 1)
216+
if (lastCodeUnit >= 0xd800 && lastCodeUnit <= 0xdbff) smallerEnd--
217+
if (smallerEnd <= offset) throw error
218+
end = smallerEnd
219+
}
220+
}
221+
}
196222
}
197223
}
198224

0 commit comments

Comments
 (0)