Skip to content
Merged
4 changes: 4 additions & 0 deletions gems/aws-sdk-s3/CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
Unreleased Changes
------------------

* Issue - Fix unbounded memory growth in `upload_stream` on `TransferManager` and `Aws::S3::Object` when the data source outpaces the upload (#3407).

* Issue - Prevent `upload_file`, `upload_stream` and `download_file` on `TransferManager` and `Aws::S3::Object` from hanging when a worker thread is terminated by a non-`StandardError`.

1.233.0 (2026-09-30)
------------------

Expand Down
7 changes: 5 additions & 2 deletions gems/aws-sdk-s3/lib/aws-sdk-s3/customizations/object.rb
Original file line number Diff line number Diff line change
Expand Up @@ -380,7 +380,9 @@ def public_url(options = {})
# and {Client#upload_part} can be provided.
#
# @option options [Integer] :thread_count (10) The number of parallel multipart uploads.
# An additional thread is used internally for task coordination.
# An additional thread is used internally for task coordination. This also bounds
# how many parts are buffered ahead of the upload, limiting memory usage to roughly
# `2 * :thread_count * :part_size`.
#
# @option options [Boolean] :tempfile (false) Normally read data is stored
# in memory when building the parts in order to complete the underlying
Expand All @@ -405,7 +407,8 @@ def public_url(options = {})
# @see Client#upload_part
def upload_stream(options = {}, &block)
upload_opts = options.merge(bucket: bucket_name, key: key)
executor = DefaultExecutor.new(max_threads: upload_opts.delete(:thread_count))
thread_count = upload_opts.delete(:thread_count) || DefaultExecutor::DEFAULT_MAX_THREADS
Comment thread
jterapin marked this conversation as resolved.
executor = DefaultExecutor.new(max_threads: thread_count, max_queue: thread_count)
Comment thread
jterapin marked this conversation as resolved.
begin
uploader = MultipartStreamUploader.new(
client: client,
Expand Down
84 changes: 65 additions & 19 deletions gems/aws-sdk-s3/lib/aws-sdk-s3/default_executor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -4,17 +4,26 @@ module Aws
module S3
# @api private
class DefaultExecutor
# Raised when a task is posted to an executor that is shutting down or has been shut down.
class RejectedExecutionError < RuntimeError
def initialize(msg = 'Executor has been shutdown and is no longer accepting tasks')
super
end
end

DEFAULT_MAX_THREADS = 10
RUNNING = :running
SHUTTING_DOWN = :shutting_down
SHUTDOWN = :shutdown

def initialize(options = {})
@max_threads = options[:max_threads] || DEFAULT_MAX_THREADS
@max_queue = options[:max_queue] || 0
@state = RUNNING
@queue = Queue.new
@queue = @max_queue.zero? ? Queue.new : SizedQueue.new(@max_queue)
@pool = []
@mutex = Mutex.new
@fatal_error = nil
end

# Submits a task for execution.
Expand All @@ -23,12 +32,15 @@ def initialize(options = {})
# @return [Boolean] Returns true if the task was submitted successfully
def post(*args, &block)
@mutex.synchronize do
raise 'Executor has been shutdown and is no longer accepting tasks' unless @state == RUNNING
raise RejectedExecutionError unless @state == RUNNING

@queue << [args, block]
ensure_worker_available
end
# Outside the mutex so a caller blocked on a full queue can't hold up #shutdown or #kill.
@queue.push([args, block])
Comment thread
jterapin marked this conversation as resolved.
true
rescue ClosedQueueError
raise RejectedExecutionError
end

# Immediately terminates all worker threads and clears pending tasks.
Expand All @@ -38,6 +50,7 @@ def post(*args, &block)
def kill
@mutex.synchronize do
@state = SHUTDOWN
@queue.close
@pool.each(&:kill)
@pool.clear
@queue.clear
Expand All @@ -52,35 +65,58 @@ def kill
# If nil, waits indefinitely. If timeout expires, remaining threads are killed.
# @return [Boolean] true when shutdown is complete
def shutdown(timeout = nil)
return true unless begin_shutdown

deadline = timeout && (Time.now + timeout)
join_workers(deadline)
kill_remaining_workers if timeout

finalize_shutdown
raise @fatal_error if @fatal_error

true
end

private

def begin_shutdown
@mutex.synchronize do
return true if @state == SHUTDOWN
return false if @state == SHUTDOWN

@state = SHUTTING_DOWN
@pool.size.times { @queue << :shutdown }
# Close rather than push sentinels, which could block on a full queue.
@queue.close
end
true
end

if timeout
deadline = Time.now + timeout
@pool.each do |thread|
remaining = deadline - Time.now
break if remaining <= 0
def join_workers(deadline)
until (threads = @mutex.synchronize { @pool.select(&:alive?) }).empty?
threads.each do |thread|
remaining = deadline - Time.now if deadline
break if remaining && remaining <= 0

thread.join([remaining, 0].max)
begin
thread.join(remaining && [remaining, 0].max)
rescue Exception # rubocop:disable Lint/RescueException
nil # recorded by #replace_worker and raised by #shutdown
end
end
@pool.select(&:alive?).each(&:kill)
else
@pool.each(&:join)
break if deadline && Time.now >= deadline
end
end

def kill_remaining_workers
@mutex.synchronize { @pool.select(&:alive?).each(&:kill) }
end

def finalize_shutdown
@mutex.synchronize do
@pool.clear
@state = SHUTDOWN
end
true
end

private

def ensure_worker_available
return unless @state == RUNNING

Expand All @@ -91,11 +127,21 @@ def ensure_worker_available
def spawn_worker
Thread.new do
while (job = @queue.shift)
break if job == :shutdown

args, block = job
block.call(*args)
Comment thread
jterapin marked this conversation as resolved.
end
rescue Exception => e # rubocop:disable Lint/RescueException
# Replace this worker so queued tasks still run.
replace_worker(e)
raise
end
end

def replace_worker(error)
@mutex.synchronize do
@fatal_error ||= error
@pool.delete(Thread.current)
@pool << spawn_worker unless @state == SHUTDOWN
end
end
end
Expand Down
36 changes: 22 additions & 14 deletions gems/aws-sdk-s3/lib/aws-sdk-s3/multipart_file_uploader.rb
Original file line number Diff line number Diff line change
Expand Up @@ -140,7 +140,7 @@ def upload_part_opts(options)
end

def upload_with_executor(pending, completed, options)
upload_attempts = 0
queued_parts = 0
completion_queue = Queue.new
abort_upload = false
errors = []
Expand All @@ -149,25 +149,33 @@ def upload_with_executor(pending, completed, options)
while (part = pending.shift)
break if abort_upload

upload_attempts += 1
@executor.post(part) do |p|
Thread.current[:net_http_override_body_stream_chunk] = @http_chunk_size if @http_chunk_size
update_progress(progress, p)
resp = @client.upload_part(p)
completed_part = { etag: resp.etag, part_number: p[:part_number] }
apply_part_checksum(resp, completed_part)
completed.push(completed_part)
begin
Comment thread
jterapin marked this conversation as resolved.
@executor.post(part) do |p|
Thread.current[:net_http_override_body_stream_chunk] = @http_chunk_size if @http_chunk_size
update_progress(progress, p)
resp = @client.upload_part(p)
completed_part = { etag: resp.etag, part_number: p[:part_number] }
apply_part_checksum(resp, completed_part)
completed.push(completed_part)
# Any error, or the upload completes without this part.
rescue Exception => e # rubocop:disable Lint/RescueException
abort_upload = true
errors << e
ensure
p[:body].close
Thread.current[:net_http_override_body_stream_chunk] = nil if @http_chunk_size
completion_queue << :done
end
queued_parts += 1
rescue StandardError => e
# Rejected by the executor; abort rather than orphan the upload.
abort_upload = true
errors << e
ensure
p[:body].close
Thread.current[:net_http_override_body_stream_chunk] = nil if @http_chunk_size
completion_queue << :done
break
end
end

upload_attempts.times { completion_queue.pop }
queued_parts.times { completion_queue.pop }
errors
end

Expand Down
49 changes: 32 additions & 17 deletions gems/aws-sdk-s3/lib/aws-sdk-s3/multipart_stream_uploader.rb
Original file line number Diff line number Diff line change
Expand Up @@ -113,18 +113,22 @@ def complete_opts(options)
def read_to_part_body(read_pipe)
return if read_pipe.closed?

temp_io = @tempfile ? Tempfile.new('aws-sdk-s3-upload_stream') : StringIO.new(String.new)
temp_io.binmode
bytes_copied = IO.copy_stream(read_pipe, temp_io, @part_size)
temp_io.rewind
if bytes_copied.zero?
if temp_io.is_a?(Tempfile)
if @tempfile
Comment thread
jterapin marked this conversation as resolved.
temp_io = Tempfile.new('aws-sdk-s3-upload_stream')
Comment thread
jterapin marked this conversation as resolved.
temp_io.binmode
bytes_copied = IO.copy_stream(read_pipe, temp_io, @part_size)
temp_io.rewind
if bytes_copied.zero?
temp_io.close
temp_io.unlink
nil
else
temp_io
end
nil
else
temp_io
# A single sized read; copy_stream into a StringIO grows by doubling and fragments the heap.
data = read_pipe.read(@part_size)
data.nil? ? nil : StringIO.new(data)
end
end

Expand All @@ -139,20 +143,31 @@ def upload_with_executor(read_pipe, completed, errors, options)
end
break unless part_body || current_part_num == 1

queued_parts += 1
@executor.post(part_body, current_part_num, options) do |body, num, opts|
part = opts.merge(body: body, part_number: num)
resp = @client.upload_part(part)
completed_part = create_completed_part(resp, part)
completed.push(completed_part)
begin
@executor.post(part_body, current_part_num, options) do |body, num, opts|
part = opts.merge(body: body, part_number: num)
resp = @client.upload_part(part)
completed_part = create_completed_part(resp, part)
completed.push(completed_part)
# Any error, or the upload completes without this part.
rescue Exception => e # rubocop:disable Lint/RescueException
mutex.synchronize do
errors.push(e)
read_pipe.close_read unless read_pipe.closed?
end
ensure
clear_body(body)
completion_queue << :done
end
queued_parts += 1
rescue StandardError => e
# Rejected by the executor. Closing the pipe stops the producer so the upload can abort.
mutex.synchronize do
errors.push(e)
read_pipe.close_read unless read_pipe.closed?
end
ensure
clear_body(body)
completion_queue << :done
clear_body(part_body)
break
end
end
queued_parts.times { completion_queue.pop }
Expand Down
7 changes: 5 additions & 2 deletions gems/aws-sdk-s3/lib/aws-sdk-s3/transfer_manager.rb
Original file line number Diff line number Diff line change
Expand Up @@ -504,6 +504,9 @@ def upload_file(source, bucket:, key:, **options)
# @option options [Integer] :thread_count (10)
# The number of parallel multipart uploads. Only used when no custom executor is provided (creates
# {DefaultExecutor} with the given thread count). An additional thread is used internally for task coordination.
# This also bounds how many parts are buffered ahead of the upload, limiting memory usage to roughly
# `2 * :thread_count * :part_size`. When a custom `:executor` is provided, it is responsible for applying
# its own backpressure.
#
# @option options [Boolean] :tempfile (false)
# Normally read data is stored in memory when building the parts in order to complete the underlying
Expand All @@ -524,8 +527,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)
thread_count = upload_opts.delete(:thread_count)
executor = @executor || DefaultExecutor.new(max_threads: thread_count)
thread_count = upload_opts.delete(:thread_count) || DefaultExecutor::DEFAULT_MAX_THREADS
executor = @executor || DefaultExecutor.new(max_threads: thread_count, max_queue: thread_count)
begin
uploader = MultipartStreamUploader.new(
client: @client,
Expand Down
Loading
Loading