From 481fbdf186dc8b166c164592c82ec9b7a5edb8b0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Iv=C3=A1n=20Ovejero?= Date: Fri, 24 Apr 2026 16:51:44 +0200 Subject: [PATCH] fix(core): Prevent S3 socket pool exhaustion on partial stream reads (#28313) --- .../__tests__/object-store.service.test.ts | 63 +++++++++++++++---- .../object-store/object-store.service.ee.ts | 43 +++++++++++-- 2 files changed, 89 insertions(+), 17 deletions(-) diff --git a/packages/core/src/binary-data/object-store/__tests__/object-store.service.test.ts b/packages/core/src/binary-data/object-store/__tests__/object-store.service.test.ts index 13155bcc217..e2b00102863 100644 --- a/packages/core/src/binary-data/object-store/__tests__/object-store.service.test.ts +++ b/packages/core/src/binary-data/object-store/__tests__/object-store.service.test.ts @@ -10,7 +10,7 @@ import { type S3Client, } from '@aws-sdk/client-s3'; import { captor, mock } from 'jest-mock-extended'; -import { Readable } from 'stream'; +import { PassThrough, Readable } from 'stream'; import type { ObjectStoreConfig } from '../object-store.config'; import { ObjectStoreService } from '../object-store.service.ee'; @@ -243,9 +243,7 @@ describe('ObjectStoreService', () => { const result = await objectStoreService.get(fileId, { mode: 'buffer' }); - const commandCaptor = captor(); - expect(mockS3Send).toHaveBeenCalledWith(commandCaptor); - const command = commandCaptor.value; + const command = mockS3Send.mock.calls[0][0] as GetObjectCommand; expect(command).toBeInstanceOf(GetObjectCommand); expect(command.input).toEqual({ Bucket: 'test-bucket', @@ -256,23 +254,66 @@ describe('ObjectStoreService', () => { }); it('should send a GET request to download an object as a stream', async () => { - const body = new Readable(); + const body = new Readable({ read() {} }); mockS3Send.mockResolvedValueOnce({ Body: body }); const result = await objectStoreService.get(fileId, { mode: 'stream' }); - const commandCaptor = captor(); - expect(mockS3Send).toHaveBeenCalledWith(commandCaptor); - const command = commandCaptor.value; - expect(command).toBeInstanceOf(GetObjectCommand); + expect(mockS3Send).toHaveBeenCalledWith( + expect.any(GetObjectCommand), + expect.objectContaining({ abortSignal: expect.any(AbortSignal) }), + ); + const command = mockS3Send.mock.calls[0][0] as GetObjectCommand; expect(command.input).toEqual({ Bucket: 'test-bucket', Key: fileId, }); - expect(result instanceof Readable).toBe(true); - expect(result).toBe(body); + expect(result).toBeInstanceOf(PassThrough); + }); + + it('should abort the S3 request when wrapper stream is destroyed before body is fully consumed', async () => { + const body = new Readable({ read() {} }); + + mockS3Send.mockResolvedValueOnce({ Body: body }); + + const result = await objectStoreService.get(fileId, { mode: 'stream' }); + + const abortSignal = mockS3Send.mock.calls[0][1].abortSignal as AbortSignal; + expect(abortSignal.aborted).toBe(false); + + result.destroy(); + + await new Promise((resolve) => result.on('close', resolve)); + + expect(abortSignal.aborted).toBe(true); + }); + + it('should pass through all data when the stream is fully consumed', async () => { + const data = 'hello world'; + let pushCount = 0; + const body = new Readable({ + read() { + if (pushCount === 0) { + this.push(data); + pushCount++; + } else { + this.push(null); + } + }, + }); + + mockS3Send.mockResolvedValueOnce({ Body: body }); + + const result = await objectStoreService.get(fileId, { mode: 'stream' }); + + const chunks: Buffer[] = []; + for await (const chunk of result) { + chunks.push(Buffer.from(chunk)); + } + + expect(Buffer.concat(chunks).toString()).toBe(data); }); it('should throw an error on request failure', async () => { diff --git a/packages/core/src/binary-data/object-store/object-store.service.ee.ts b/packages/core/src/binary-data/object-store/object-store.service.ee.ts index 91b7317822f..ab2ffe62dd4 100644 --- a/packages/core/src/binary-data/object-store/object-store.service.ee.ts +++ b/packages/core/src/binary-data/object-store/object-store.service.ee.ts @@ -18,7 +18,7 @@ import { Logger } from '@n8n/backend-common'; import { Service } from '@n8n/di'; import { UnexpectedError } from 'n8n-workflow'; import { createHash } from 'node:crypto'; -import { Readable } from 'node:stream'; +import { PassThrough, Readable } from 'node:stream'; import { ObjectStoreConfig } from './object-store.config'; import type { MetadataResponseHeaders } from './types'; @@ -137,14 +137,45 @@ export class ObjectStoreService { }); try { + if (mode === 'stream') { + const abortController = new AbortController(); + const { Body: body } = await this.s3Client.send(command, { + abortSignal: abortController.signal, + }); + if (!body) throw new UnexpectedError('Received empty response body'); + + if (!(body instanceof Readable)) { + throw new UnexpectedError('Expected stream but received different type', { + extra: { bodyType: typeof body }, + }); + } + + // Wrap to prevent socket pool exhaustion when callers destroy the + // stream early. AbortController lets the SDK free the socket slot + // properly. See: https://github.com/aws/aws-sdk-js-v3/issues/6691 + const wrapper = new PassThrough(); + let bodyFullyConsumed = false; + + body.on('end', () => { + bodyFullyConsumed = true; + }); + + wrapper.on('close', () => { + if (!bodyFullyConsumed) { + abortController.abort(); + body.destroy(); + } + }); + + body.on('error', (error) => wrapper.destroy(error)); + body.pipe(wrapper); + + return wrapper; + } + const { Body: body } = await this.s3Client.send(command); if (!body) throw new UnexpectedError('Received empty response body'); - if (mode === 'stream') { - if (body instanceof Readable) return body; - throw new UnexpectedError(`Expected stream but received ${typeof body}.`); - } - return await streamToBuffer(body as Readable); } catch (e) { throw new UnexpectedError('Request to S3 failed', { cause: e });