Skip to content

Commit 72eb8b6

Browse files
committed
fix(logs): guard export stream cancellation
1 parent 3ca60d2 commit 72eb8b6

2 files changed

Lines changed: 26 additions & 0 deletions

File tree

apps/sim/app/api/logs/export/route.test.ts

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -214,4 +214,26 @@ describe('GET /api/logs/export', () => {
214214
await reader.cancel()
215215
expect(dbChainMockFns.where).toHaveBeenCalledTimes(1)
216216
})
217+
218+
it('stops a pending pull cleanly when the reader cancels', async () => {
219+
queueTableRows(workflowExecutionLogs, [logRow(0)])
220+
let resolveMaterialization: ((value: unknown[]) => void) | undefined
221+
mockMapWithConcurrency.mockImplementationOnce(
222+
() =>
223+
new Promise((resolve) => {
224+
resolveMaterialization = resolve
225+
})
226+
)
227+
228+
const response = await GET(makeRequest())
229+
const reader = response.body!.getReader()
230+
await reader.read()
231+
232+
const pendingRead = reader.read()
233+
await vi.waitFor(() => expect(mockMapWithConcurrency).toHaveBeenCalledTimes(1))
234+
const cancellation = reader.cancel()
235+
resolveMaterialization?.([{ message: 'message-0' }])
236+
237+
await expect(Promise.all([pendingRead, cancellation])).resolves.toBeDefined()
238+
})
217239
})

apps/sim/app/api/logs/export/route.ts

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -177,22 +177,26 @@ export const GET = withRouteHandler(async (request: NextRequest) => {
177177
}
178178
})()
179179

180+
let cancelled = false
180181
const stream = new ReadableStream<Uint8Array>(
181182
{
182183
pull: async (controller) => {
183184
try {
184185
const next = await csvChunks.next()
186+
if (cancelled) return
185187
if (next.done) {
186188
controller.close()
187189
return
188190
}
189191
controller.enqueue(next.value)
190192
} catch (error) {
193+
if (cancelled) return
191194
logger.error('Export stream error', { error: getErrorMessage(error) })
192195
controller.error(error)
193196
}
194197
},
195198
cancel: async () => {
199+
cancelled = true
196200
await csvChunks.return(undefined)
197201
},
198202
},

0 commit comments

Comments
 (0)