Skip to content

Commit 4a4ca0b

Browse files
fix(slack): preserve text order around tool progress
1 parent dababea commit 4a4ca0b

2 files changed

Lines changed: 106 additions & 5 deletions

File tree

‎apps/sim/lib/slack-search/assistant-stream.test.ts‎

Lines changed: 95 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,101 @@ function toolResult(
111111
}
112112

113113
describe('Slack tool progress', () => {
114+
it('flushes a batched sentence before starting tool progress', async () => {
115+
vi.spyOn(Date, 'now').mockReturnValue(1000)
116+
try {
117+
const { stream } = setup()
118+
await stream.start()
119+
await stream.onEvent({
120+
type: 'text',
121+
payload: { channel: 'assistant', text: "I'll search " },
122+
})
123+
await stream.onEvent({
124+
type: 'text',
125+
payload: { channel: 'assistant', text: 'the connected sources for the handbook.' },
126+
})
127+
await stream.onEvent(toolCall('search_workspace'))
128+
const chunks = api.append.mock.calls.flatMap((call) => call[3])
129+
expect(chunks).toEqual([
130+
{ type: 'markdown_text', text: "I'll search " },
131+
{
132+
type: 'markdown_text',
133+
text: 'the connected sources for the handbook.\n\n',
134+
},
135+
{
136+
type: 'task_update',
137+
id: expect.any(String),
138+
title: 'Searching documents…',
139+
status: 'in_progress',
140+
},
141+
])
142+
} finally {
143+
vi.restoreAllMocks()
144+
}
145+
})
146+
147+
it('keeps text contiguous across preparatory and hidden tool events', async () => {
148+
const { stream } = setup()
149+
await stream.start()
150+
await stream.onEvent({
151+
type: 'text',
152+
payload: { channel: 'assistant', text: "I'll search" },
153+
})
154+
for (const attributes of [
155+
{ partial: true },
156+
{ status: 'generating' as const },
157+
{ ui: { hidden: true } },
158+
{ ui: { internal: true } },
159+
]) {
160+
const event = toolCall('search_workspace')
161+
await stream.onEvent({ ...event, payload: { ...event.payload, ...attributes } })
162+
}
163+
await stream.onEvent({
164+
type: 'text',
165+
payload: { channel: 'assistant', text: ' the connected sources.' },
166+
})
167+
await stream.finish(result)
168+
expect(deliveredText()).toBe("I'll search the connected sources.")
169+
expect(
170+
api.append.mock.calls
171+
.flatMap((call) => call[3])
172+
.every((chunk) => chunk.type === 'markdown_text')
173+
).toBe(true)
174+
})
175+
176+
it('serializes concurrent text and tool events without duplicating buffered text', async () => {
177+
const { stream } = setup()
178+
await stream.start()
179+
let releaseAppend!: () => void
180+
api.append.mockImplementationOnce(
181+
() =>
182+
new Promise<void>((resolve) => {
183+
releaseAppend = resolve
184+
})
185+
)
186+
const text = stream.onEvent({
187+
type: 'text',
188+
payload: { channel: 'assistant', text: "I'll search the connected sources. " },
189+
})
190+
await vi.waitFor(() => expect(api.append).toHaveBeenCalledOnce(), { interval: 1 })
191+
const call = stream.onEvent(toolCall('search_workspace'))
192+
const completed = stream.onEvent(toolResult('search_workspace'))
193+
const finished = stream.finish(result)
194+
expect(api.append).toHaveBeenCalledOnce()
195+
expect(api.stop).not.toHaveBeenCalled()
196+
releaseAppend()
197+
await Promise.all([text, call, completed, finished])
198+
const chunks = api.append.mock.calls.flatMap((call) => call[3])
199+
expect(deliveredText()).toBe("I'll search the connected sources. \n\n")
200+
expect(chunks.map((chunk) => chunk.type)).toEqual([
201+
'markdown_text',
202+
'markdown_text',
203+
'task_update',
204+
'task_update',
205+
])
206+
expect(chunks[3]).toEqual({ ...chunks[2], status: 'complete' })
207+
})
208+
114209
it.each([
115210
['list_integrations', 'Listing connected integrations…'],
116211
['search_workspace', 'Searching documents…'],

‎apps/sim/lib/slack-search/assistant-stream.ts‎

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -123,7 +123,7 @@ export class SlackSearchAssistantStream {
123123
private failure?: Error
124124
private closed = false
125125
private closeAttempted = false
126-
private separateNextText = false
126+
private pendingEvents: Promise<void> = Promise.resolve()
127127
private evidence = new Map<string, Record<string, unknown>>()
128128
private toolProgress = new Map<string, { toolName: string; chunk: ToolProgress }>()
129129
constructor(private readonly options: AssistantStreamOptions) {}
@@ -160,7 +160,12 @@ export class SlackSearchAssistantStream {
160160
})
161161
}
162162

163-
async onEvent(event: StreamEvent) {
163+
onEvent(event: StreamEvent): Promise<void> {
164+
this.pendingEvents = this.pendingEvents.then(() => this.handleEvent(event))
165+
return this.pendingEvents
166+
}
167+
168+
private async handleEvent(event: StreamEvent) {
164169
if (this.failure) throw this.failure
165170
if (event.type === 'tool' && 'phase' in event.payload && event.payload.phase === 'result') {
166171
const { toolName, success, status, output } = event.payload
@@ -175,7 +180,6 @@ export class SlackSearchAssistantStream {
175180
])
176181
}
177182
if (event.type === 'tool' && !event.scope) {
178-
this.separateNextText = true
179183
if (
180184
'phase' in event.payload &&
181185
(event.payload.phase === 'call' || event.payload.phase === 'result')
@@ -184,8 +188,6 @@ export class SlackSearchAssistantStream {
184188
}
185189
}
186190
if (event.type !== 'text' || event.payload.channel !== 'assistant' || event.scope) return
187-
if (this.separateNextText && this.text) this.text += '\n\n'
188-
this.separateNextText = false
189191
this.text += event.payload.text
190192
if (this.text.length > 128_000) throw new Error('Slack answer exceeds the supported size')
191193
if (Date.now() - this.lastSentAt >= 750) await this.flush(false)
@@ -208,6 +210,9 @@ export class SlackSearchAssistantStream {
208210
(payload.status !== undefined && payload.status !== 'executing')
209211
)
210212
return
213+
/** Close the preceding text segment so batching cannot place its tail after the task. */
214+
if (this.text && !this.text.endsWith('\n\n')) this.text += '\n\n'
215+
await this.flush(false)
211216
chunk = { type: 'task_update', id: generateId(), title, status: 'in_progress' }
212217
this.toolProgress.set(payload.toolCallId, { toolName: payload.toolName, chunk })
213218
} else {
@@ -298,6 +303,7 @@ export class SlackSearchAssistantStream {
298303
}
299304

300305
async finish(result: OrchestratorResult) {
306+
await this.pendingEvents
301307
if (this.failure) throw this.failure
302308
this.collectSources(result.contentBlocks)
303309
const projection = projectResolvedSecretDiagnosticContent(

0 commit comments

Comments
 (0)