mirror of
https://github.com/n8n-io/n8n.git
synced 2026-09-24 23:22:38 +08:00
fix(core): Prevent S3 socket pool exhaustion on partial stream reads (#28313)
This commit is contained in:
@@ -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<GetObjectCommand>();
|
||||
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<GetObjectCommand>();
|
||||
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 () => {
|
||||
|
||||
@@ -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 });
|
||||
|
||||
Reference in New Issue
Block a user