From 6477135deb7ebd55dff3edf512ba52a2d45af1d1 Mon Sep 17 00:00:00 2001 From: Chaitanya Bhorade Date: Wed, 23 Sep 2026 16:59:35 -0700 Subject: [PATCH 1/3] Shutdown internally-created executor on error in TransferManager --- gems/aws-sdk-s3/CHANGELOG.md | 2 + .../lib/aws-sdk-s3/transfer_manager.rb | 75 +++++++++++-------- gems/aws-sdk-s3/spec/transfer_manager_spec.rb | 74 ++++++++++++++++++ 3 files changed, 120 insertions(+), 31 deletions(-) diff --git a/gems/aws-sdk-s3/CHANGELOG.md b/gems/aws-sdk-s3/CHANGELOG.md index 129db9dbd12..a6472af8139 100644 --- a/gems/aws-sdk-s3/CHANGELOG.md +++ b/gems/aws-sdk-s3/CHANGELOG.md @@ -1,6 +1,8 @@ Unreleased Changes ------------------ +* Issue - Ensure the internally-created executor is shutdown on error in `TransferManager` methods, preventing leaked worker threads on multipart transfer failures (#3419). + 1.232.1 (2026-09-16) ------------------ diff --git a/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb b/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb index 864dd543c14..7668bfdde0e 100644 --- a/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb +++ b/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb @@ -161,10 +161,12 @@ def initialize(options = {}) def download_directory(destination, bucket:, **options) Aws::Plugins::UserAgent.metric('S3_TRANSFER', 'S3_TRANSFER_DOWNLOAD_DIRECTORY') do executor = @executor || DefaultExecutor.new(max_threads: options.delete(:thread_count)) - downloader = DirectoryDownloader.new(client: @client, executor: executor, logger: @logger) - result = downloader.download(destination, bucket: bucket, **options) - executor.shutdown unless @executor - result + begin + downloader = DirectoryDownloader.new(client: @client, executor: executor, logger: @logger) + downloader.download(destination, bucket: bucket, **options) + ensure + executor.shutdown unless @executor + end end end @@ -246,10 +248,13 @@ def download_directory(destination, bucket:, **options) def download_file(destination, bucket:, key:, **options) download_opts = options.merge(bucket: bucket, key: key) executor = @executor || DefaultExecutor.new(max_threads: download_opts.delete(:thread_count)) - downloader = FileDownloader.new(client: @client, executor: executor) - downloader.download(destination, download_opts) - executor.shutdown unless @executor - true + begin + downloader = FileDownloader.new(client: @client, executor: executor) + downloader.download(destination, download_opts) + true + ensure + executor.shutdown unless @executor + end end # Uploads all files under the given directory to the provided S3 bucket. @@ -361,10 +366,12 @@ def download_file(destination, bucket:, key:, **options) def upload_directory(source, bucket:, **options) Aws::Plugins::UserAgent.metric('S3_TRANSFER', 'S3_TRANSFER_UPLOAD_DIRECTORY') do executor = @executor || DefaultExecutor.new(max_threads: options.delete(:thread_count)) - uploader = DirectoryUploader.new(client: @client, executor: executor, logger: @logger) - result = uploader.upload(source, bucket, **options.merge(http_chunk_size: resolve_http_chunk_size(options))) - executor.shutdown unless @executor - result + begin + uploader = DirectoryUploader.new(client: @client, executor: executor, logger: @logger) + uploader.upload(source, bucket, **options.merge(http_chunk_size: resolve_http_chunk_size(options))) + ensure + executor.shutdown unless @executor + end end end @@ -445,16 +452,19 @@ def upload_file(source, bucket:, key:, **options) http_chunk_size = resolve_http_chunk_size(upload_opts) executor = @executor || DefaultExecutor.new(max_threads: upload_opts.delete(:thread_count)) - uploader = FileUploader.new( - multipart_threshold: upload_opts.delete(:multipart_threshold), - http_chunk_size: http_chunk_size, - client: @client, - executor: executor - ) - response = uploader.upload(source, upload_opts) - yield response if block_given? - executor.shutdown unless @executor - true + begin + uploader = FileUploader.new( + multipart_threshold: upload_opts.delete(:multipart_threshold), + http_chunk_size: http_chunk_size, + client: @client, + executor: executor + ) + response = uploader.upload(source, upload_opts) + yield response if block_given? + true + ensure + executor.shutdown unless @executor + end end # Uploads a stream in a streaming fashion to S3. @@ -512,15 +522,18 @@ def upload_file(source, bucket:, key:, **options) def upload_stream(bucket:, key:, **options, &block) upload_opts = options.merge(bucket: bucket, key: key) executor = @executor || DefaultExecutor.new(max_threads: upload_opts.delete(:thread_count)) - uploader = MultipartStreamUploader.new( - client: @client, - executor: executor, - tempfile: upload_opts.delete(:tempfile), - part_size: upload_opts.delete(:part_size) - ) - uploader.upload(upload_opts, &block) - executor.shutdown unless @executor - true + begin + uploader = MultipartStreamUploader.new( + client: @client, + executor: executor, + tempfile: upload_opts.delete(:tempfile), + part_size: upload_opts.delete(:part_size) + ) + uploader.upload(upload_opts, &block) + true + ensure + executor.shutdown unless @executor + end end private diff --git a/gems/aws-sdk-s3/spec/transfer_manager_spec.rb b/gems/aws-sdk-s3/spec/transfer_manager_spec.rb index d938428e2c2..036fb3725e9 100644 --- a/gems/aws-sdk-s3/spec/transfer_manager_spec.rb +++ b/gems/aws-sdk-s3/spec/transfer_manager_spec.rb @@ -51,6 +51,23 @@ module S3 subject.download_directory(temp_dir, bucket: 'bucket', ignore_failure: false) end.to raise_error(DirectoryDownloadError) end + + it 'shuts down the internally-created executor when download raises' do + client.stub_responses(:get_object, 'AccessDenied') + executor = nil + allow(DefaultExecutor).to receive(:new).and_wrap_original do |orig, *args, **kwargs| + instance = orig.call(*args, **kwargs) + if executor.nil? + executor = instance + allow(instance).to receive(:shutdown).and_call_original + end + instance + end + expect do + subject.download_directory(temp_dir, bucket: 'bucket', ignore_failure: false) + end.to raise_error(DirectoryDownloadError) + expect(executor).to have_received(:shutdown) + end end describe '#download_file' do @@ -72,6 +89,19 @@ module S3 .to raise_error(Aws::S3::Errors::NoSuchKey) end + it 'shuts down the internally-created executor when download raises' do + client.stub_responses(:head_object, 'NoSuchKey') + executor = nil + allow(DefaultExecutor).to receive(:new).and_wrap_original do |orig, *args, **kwargs| + executor = orig.call(*args, **kwargs) + allow(executor).to receive(:shutdown).and_call_original + executor + end + expect { subject.download_file(path, bucket: 'bucket', key: 'missing-key') } + .to raise_error(Aws::S3::Errors::NoSuchKey) + expect(executor).to have_received(:shutdown) + end + it 'calls progress callback when given' do n_calls = 0 callback = proc { |_b, _p, _t| n_calls += 1 } @@ -108,6 +138,23 @@ module S3 subject.upload_directory(temp_dir, bucket: 'bucket', ignore_failure: false) end.to raise_error(DirectoryUploadError) end + + it 'shuts down the internally-created executor when upload raises' do + client.stub_responses(:put_object, 'AccessDenied') + executor = nil + allow(DefaultExecutor).to receive(:new).and_wrap_original do |orig, *args, **kwargs| + instance = orig.call(*args, **kwargs) + if executor.nil? + executor = instance + allow(instance).to receive(:shutdown).and_call_original + end + instance + end + expect do + subject.upload_directory(temp_dir, bucket: 'bucket', ignore_failure: false) + end.to raise_error(DirectoryUploadError) + expect(executor).to have_received(:shutdown) + end end describe '#upload_file' do @@ -135,6 +182,19 @@ module S3 .to raise_error(Aws::S3::Errors::AccessDenied) end + it 'shuts down the internally-created executor when upload raises' do + client.stub_responses(:put_object, 'AccessDenied') + executor = nil + allow(DefaultExecutor).to receive(:new).and_wrap_original do |orig, *args, **kwargs| + executor = orig.call(*args, **kwargs) + allow(executor).to receive(:shutdown).and_call_original + executor + end + expect { subject.upload_file(file, bucket: 'forbidden-bucket', key: 'key') } + .to raise_error(Aws::S3::Errors::AccessDenied) + expect(executor).to have_received(:shutdown) + end + it 'yields response when block given' do subject.upload_file(file, bucket: 'bucket', key: 'key') do |response| expect(response).to be_kind_of(Seahorse::Client::Response) @@ -243,6 +303,20 @@ module S3 subject.upload_stream(bucket: 'bucket', key: 'key') { |write_stream| write_stream << seventeen_mb } end.to raise_error(Aws::S3::MultipartUploadError, /part failed/) end + + it 'shuts down the internally-created executor when upload raises' do + client.stub_responses(:upload_part, RuntimeError.new('part failed')) + executor = nil + allow(DefaultExecutor).to receive(:new).and_wrap_original do |orig, *args, **kwargs| + executor = orig.call(*args, **kwargs) + allow(executor).to receive(:shutdown).and_call_original + executor + end + expect do + subject.upload_stream(bucket: 'bucket', key: 'key') { |write_stream| write_stream << seventeen_mb } + end.to raise_error(Aws::S3::MultipartUploadError, /part failed/) + expect(executor).to have_received(:shutdown) + end end end end From c2a87f651c815ac82baa076fc18ac32c3df52dae Mon Sep 17 00:00:00 2001 From: Juli Tera Date: Thu, 24 Sep 2026 10:15:48 -0700 Subject: [PATCH 2/3] Fix :thread_count handling in TransferManager and Object transfer methods The five TransferManager transfer methods build their executor as `@executor || DefaultExecutor.new(max_threads: opts.delete(:thread_count))`. Ruby's `||` short-circuits, so when a custom executor was given to the constructor the `delete` never runs and `:thread_count` stays in the options hash. On `upload_file` it reaches the unfiltered `put_object` call on the single-part path and raises `ArgumentError: unexpected value at params[:thread_count]`, while the multipart path filters the key out, so the same call fails or succeeds depending on file size. Hoist the `delete` above the `||` so the key is always removed. Behavior is now uniform: `:thread_count` is honored when no custom executor was provided and ignored when one was, matching the existing docstrings. Separately, `Object#download_file` deleted `[:thread_count]`, an Array that never matches the Symbol, so it always used the default thread count. --- gems/aws-sdk-s3/CHANGELOG.md | 4 + .../lib/aws-sdk-s3/customizations/object.rb | 2 +- .../lib/aws-sdk-s3/transfer_manager.rb | 15 ++- .../spec/object/download_file_spec.rb | 6 + gems/aws-sdk-s3/spec/transfer_manager_spec.rb | 125 ++++++++++++++++++ 5 files changed, 146 insertions(+), 6 deletions(-) diff --git a/gems/aws-sdk-s3/CHANGELOG.md b/gems/aws-sdk-s3/CHANGELOG.md index a6472af8139..06f4aa77aff 100644 --- a/gems/aws-sdk-s3/CHANGELOG.md +++ b/gems/aws-sdk-s3/CHANGELOG.md @@ -3,6 +3,10 @@ Unreleased Changes * Issue - Ensure the internally-created executor is shutdown on error in `TransferManager` methods, preventing leaked worker threads on multipart transfer failures (#3419). +* Issue - Remove `:thread_count` from transfer options even when a custom `:executor` is configured on `TransferManager`, fixing an `ArgumentError` raised on single-part `upload_file` calls. + +* Issue - Fix `:thread_count` being ignored by the deprecated `Object#download_file`, which always used the default thread count. + 1.232.1 (2026-09-16) ------------------ diff --git a/gems/aws-sdk-s3/lib/aws-sdk-s3/customizations/object.rb b/gems/aws-sdk-s3/lib/aws-sdk-s3/customizations/object.rb index dd2cb036edf..f0a8cac2ef4 100644 --- a/gems/aws-sdk-s3/lib/aws-sdk-s3/customizations/object.rb +++ b/gems/aws-sdk-s3/lib/aws-sdk-s3/customizations/object.rb @@ -561,7 +561,7 @@ def upload_file(source, options = {}) # @see Client#head_object def download_file(destination, options = {}) download_opts = options.merge(bucket: bucket_name, key: key) - executor = DefaultExecutor.new(max_threads: download_opts.delete([:thread_count])) + executor = DefaultExecutor.new(max_threads: download_opts.delete(:thread_count)) downloader = FileDownloader.new(client: client, executor: executor) Aws::Plugins::UserAgent.metric('RESOURCE_MODEL') do downloader.download(destination, download_opts) diff --git a/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb b/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb index 7668bfdde0e..e0507d6986d 100644 --- a/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb +++ b/gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb @@ -160,7 +160,8 @@ def initialize(options = {}) # * `:errors` - Array of errors for failed downloads (only present when failures occur) def download_directory(destination, bucket:, **options) Aws::Plugins::UserAgent.metric('S3_TRANSFER', 'S3_TRANSFER_DOWNLOAD_DIRECTORY') do - executor = @executor || DefaultExecutor.new(max_threads: options.delete(:thread_count)) + thread_count = options.delete(:thread_count) + executor = @executor || DefaultExecutor.new(max_threads: thread_count) begin downloader = DirectoryDownloader.new(client: @client, executor: executor, logger: @logger) downloader.download(destination, bucket: bucket, **options) @@ -247,7 +248,8 @@ def download_directory(destination, bucket:, **options) # @see Client#head_object def download_file(destination, bucket:, key:, **options) download_opts = options.merge(bucket: bucket, key: key) - executor = @executor || DefaultExecutor.new(max_threads: download_opts.delete(:thread_count)) + thread_count = download_opts.delete(:thread_count) + executor = @executor || DefaultExecutor.new(max_threads: thread_count) begin downloader = FileDownloader.new(client: @client, executor: executor) downloader.download(destination, download_opts) @@ -365,7 +367,8 @@ def download_file(destination, bucket:, key:, **options) # * `:errors` - Array of error objects for failed uploads (only present when failures occur) def upload_directory(source, bucket:, **options) Aws::Plugins::UserAgent.metric('S3_TRANSFER', 'S3_TRANSFER_UPLOAD_DIRECTORY') do - executor = @executor || DefaultExecutor.new(max_threads: options.delete(:thread_count)) + thread_count = options.delete(:thread_count) + executor = @executor || DefaultExecutor.new(max_threads: thread_count) begin uploader = DirectoryUploader.new(client: @client, executor: executor, logger: @logger) uploader.upload(source, bucket, **options.merge(http_chunk_size: resolve_http_chunk_size(options))) @@ -451,7 +454,8 @@ def upload_file(source, bucket:, key:, **options) upload_opts = options.merge(bucket: bucket, key: key) http_chunk_size = resolve_http_chunk_size(upload_opts) - executor = @executor || DefaultExecutor.new(max_threads: upload_opts.delete(:thread_count)) + thread_count = upload_opts.delete(:thread_count) + executor = @executor || DefaultExecutor.new(max_threads: thread_count) begin uploader = FileUploader.new( multipart_threshold: upload_opts.delete(:multipart_threshold), @@ -521,7 +525,8 @@ def upload_file(source, bucket:, key:, **options) # @see Client#upload_part def upload_stream(bucket:, key:, **options, &block) upload_opts = options.merge(bucket: bucket, key: key) - executor = @executor || DefaultExecutor.new(max_threads: upload_opts.delete(:thread_count)) + thread_count = upload_opts.delete(:thread_count) + executor = @executor || DefaultExecutor.new(max_threads: thread_count) begin uploader = MultipartStreamUploader.new( client: @client, diff --git a/gems/aws-sdk-s3/spec/object/download_file_spec.rb b/gems/aws-sdk-s3/spec/object/download_file_spec.rb index 455a51a5113..38dec27905e 100644 --- a/gems/aws-sdk-s3/spec/object/download_file_spec.rb +++ b/gems/aws-sdk-s3/spec/object/download_file_spec.rb @@ -28,6 +28,12 @@ module S3 expect { subject.download_file(path) }.to raise_error(Aws::S3::Errors::NoSuchKey) end + it 'respects the thread_count option' do + custom_thread_count = 20 + expect(DefaultExecutor).to receive(:new).with(max_threads: custom_thread_count).and_call_original + subject.download_file(path, thread_count: custom_thread_count) + end + it 'calls progress callback when given' do n_calls = 0 callback = proc { |_b, _p, _t| n_calls += 1 } diff --git a/gems/aws-sdk-s3/spec/transfer_manager_spec.rb b/gems/aws-sdk-s3/spec/transfer_manager_spec.rb index 036fb3725e9..02282dd29a9 100644 --- a/gems/aws-sdk-s3/spec/transfer_manager_spec.rb +++ b/gems/aws-sdk-s3/spec/transfer_manager_spec.rb @@ -318,6 +318,131 @@ module S3 expect(executor).to have_received(:shutdown) end end + + describe ':thread_count' do + let(:custom_executor) { DefaultExecutor.new(max_threads: 3) } + let(:temp_dir) { Dir.mktmpdir } + let(:destination) { Tempfile.new('destination').path } + let(:small_file) { Tempfile.new('small-file').tap { |f| f.write('.' * 100) && f.rewind } } + let(:ten_mb_file) do + Tempfile.new('ten-meg-file').tap do |f| + 10.times { f.write(one_mb_content) } + f.rewind + end + end + let(:seventeen_mb) { one_mb_content * 17 } + + # Records the :max_threads of every executor the SDK builds internally. The directory + # methods also build a queue executor with max_threads: 2, so assertions check for the + # requested value rather than the size of the list. + let(:max_threads) do + [].tap do |seen| + allow(DefaultExecutor).to receive(:new).and_wrap_original do |orig, *args, **kwargs| + seen << kwargs[:max_threads] + orig.call(*args, **kwargs) + end + end + end + + before do + max_threads + client.stub_responses(:head_object, content_length: one_mb_size, parts_count: nil) + client.stub_responses(:get_object, { body: 'hello-world' }) + client.stub_responses( + :list_objects_v2, + { contents: [{ key: 'file1.txt', size: 100 }], is_truncated: false } + ) + TransferManagerSpecHelper.create_test_directory_structure(temp_dir) + end + + after do + FileUtils.rm_rf(temp_dir) + custom_executor.shutdown + end + + context 'when no executor was provided to the TransferManager' do + it 'is used for #upload_file' do + subject.upload_file(small_file, bucket: 'bucket', key: 'key', thread_count: 4) + expect(max_threads).to include(4) + end + + it 'is used for #download_file' do + subject.download_file(destination, bucket: 'bucket', key: 'key', thread_count: 4) + expect(max_threads).to include(4) + end + + it 'is used for #upload_stream' do + subject.upload_stream(bucket: 'bucket', key: 'key', thread_count: 4) do |stream| + stream << seventeen_mb + end + expect(max_threads).to include(4) + end + + it 'is used for #upload_directory' do + subject.upload_directory(temp_dir, bucket: 'bucket', thread_count: 4) + expect(max_threads).to include(4) + end + + it 'is used for #download_directory' do + subject.download_directory(Dir.mktmpdir, bucket: 'bucket', thread_count: 4) + expect(max_threads).to include(4) + end + end + + context 'when an executor was provided to the TransferManager' do + let(:subject) { TransferManager.new(client: client, executor: custom_executor) } + + it 'is ignored by #upload_file' do + expect do + subject.upload_file(small_file, bucket: 'bucket', key: 'key', thread_count: 4) + end.not_to raise_error + expect(max_threads).not_to include(4) + end + + it 'is ignored by #download_file' do + expect do + subject.download_file(destination, bucket: 'bucket', key: 'key', thread_count: 4) + end.not_to raise_error + expect(max_threads).not_to include(4) + end + + it 'is ignored by #upload_stream' do + expect do + subject.upload_stream(bucket: 'bucket', key: 'key', thread_count: 4) do |stream| + stream << seventeen_mb + end + end.not_to raise_error + expect(max_threads).not_to include(4) + end + + it 'is ignored by #upload_directory' do + expect do + subject.upload_directory(temp_dir, bucket: 'bucket', thread_count: 4) + end.not_to raise_error + expect(max_threads).not_to include(4) + end + + it 'is ignored by #download_directory' do + expect do + subject.download_directory(Dir.mktmpdir, bucket: 'bucket', thread_count: 4) + end.not_to raise_error + expect(max_threads).not_to include(4) + end + + # Regression: the option used to leak into the unfiltered #put_object call on the + # single-part path, so an identical call raised or not depending on the file size. + it 'does not leak into the request params on the single-part upload path' do + expect(client).to receive(:put_object).with(hash_excluding(:thread_count)).and_call_original + subject.upload_file(small_file, bucket: 'bucket', key: 'key', thread_count: 4) + end + + it 'does not leak into the request params on the multipart upload path' do + expect do + subject.upload_file(ten_mb_file, bucket: 'bucket', key: 'key', thread_count: 4) + end.not_to raise_error + end + end + end end end end From 783248e6b759b3cc2cba5ff5adbf2c51a2200bb1 Mon Sep 17 00:00:00 2001 From: Juli Tera Date: Mon, 28 Sep 2026 08:37:57 -0700 Subject: [PATCH 3/3] docs: restore S3 thread count changelog entries --- gems/aws-sdk-s3/CHANGELOG.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/gems/aws-sdk-s3/CHANGELOG.md b/gems/aws-sdk-s3/CHANGELOG.md index c0675a2a23c..872162f3953 100644 --- a/gems/aws-sdk-s3/CHANGELOG.md +++ b/gems/aws-sdk-s3/CHANGELOG.md @@ -1,6 +1,10 @@ Unreleased Changes ------------------ +* Issue - Remove `:thread_count` from transfer options even when a custom `:executor` is configured on `TransferManager`, fixing an `ArgumentError` raised on single-part `upload_file` calls. + +* Issue - Fix `:thread_count` being ignored by the deprecated `Object#download_file`, which always used the default thread count. + 1.232.2 (2026-09-25) ------------------