diff --git a/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener.rb b/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener.rb index 814ed067ee7d..24f1287df291 100644 --- a/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener.rb +++ b/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener.rb @@ -71,6 +71,7 @@ class MessageListener attr_reader :message_ordering attr_reader :callback_threads attr_reader :push_threads + attr_reader :shutdown_behavior ## # @private Implementation attributes. @@ -91,6 +92,7 @@ def initialize subscription_name, callback, deadline: nil, message_ordering: nil @subscription_name = subscription_name @deadline = deadline || 60 @streams = streams || 1 + @shutdown_behavior = :wait_for_processing coerce_inventory inventory @message_ordering = message_ordering @callback_threads = Integer(threads[:callback] || 8) diff --git a/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/inventory.rb b/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/inventory.rb index e0fd50fd23cc..8a557c655e80 100644 --- a/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/inventory.rb +++ b/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/inventory.rb @@ -109,6 +109,31 @@ def empty? end end + ## + # @private + # Blocks until the inventory is empty or until timeout expires. + # + # @param [Numeric, nil] timeout Maximum time in seconds to wait, or nil to wait indefinitely. + # @return [Boolean] true if inventory became empty, false if timed out. + def wait_until_empty timeout = nil + synchronize do + return true if @inventory.empty? + + if timeout + target_time = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout + while !@inventory.empty? && !@stopped + remaining = target_time - Process.clock_gettime(Process::CLOCK_MONOTONIC) + break if remaining <= 0 + + @wait_cond.wait remaining + end + else + @wait_cond.wait_while { !@inventory.empty? && !@stopped } + end + @inventory.empty? + end + end + def start @background_thread ||= Thread.new { background_run } diff --git a/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/stream.rb b/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/stream.rb index 477e8bc8a585..c6272bbb2487 100644 --- a/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/stream.rb +++ b/google-cloud-pubsub/lib/google/cloud/pubsub/message_listener/stream.rb @@ -111,6 +111,11 @@ def start self end + ## + # @private + # Stops pulling messages from the subscription. + # + # @return [Stream] self for chaining. def stop synchronize do break if @stopped @@ -129,22 +134,44 @@ def stop @keepalive_monitor.stop - # Now that the reception thread is stopped, immediately stop the - # callback thread pool. All queued callbacks will see the stream - # is stopped and perform a noop. - @callback_thread_pool.shutdown - - # Once all the callbacks are stopped, we can stop the inventory. - @inventory.stop + # When :nack_immediately, release all queued messages immediately and shut down callback pool. + if nack_immediately? + nack_unprocessed_messages! + @callback_thread_pool.shutdown + end end self end + ## + # @private + # Nacks all messages currently held in inventory. + def nack_unprocessed_messages! + synchronize do + ack_ids = @inventory.ack_ids + return if ack_ids.empty? + + @subscriber.buffer.modify_ack_deadline 0, ack_ids + @inventory.remove ack_ids + end + end + def stopped? synchronize { @stopped } end + ## + # @private + # Returns whether the subscriber is stopped and configured to nack unprocessed messages immediately. + # + # @return [Boolean] + def nack_immediately? + return false unless @stopped + + @subscriber.shutdown_behavior == :nack_immediately + end + def stream_open? synchronize { @stream_open } end @@ -157,14 +184,26 @@ def running? !stopped? end + ## + # @private + # Blocks until all received messages are processed, or until timeout expires. + # + # @param [Numeric, nil] timeout The maximum seconds to wait, or nil to wait indefinitely. + # @return [Stream] self for chaining. def wait! timeout = nil - # Wait for all queued callbacks to be processed. - @callback_thread_pool.wait_for_termination timeout + deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout if timeout + emptied = @inventory.wait_until_empty timeout + nack_unprocessed_messages! unless emptied + @callback_thread_pool.shutdown + pool_timeout = [deadline - Process.clock_gettime(Process::CLOCK_MONOTONIC), 0].max if deadline + @callback_thread_pool.wait_for_termination pool_timeout + + # Once all callbacks are finished and inventory is clear, stop the inventory. + @inventory.stop self end - def request_queue_active? !@request_queue.nil? end @@ -466,7 +505,7 @@ def perform_callback_sync rec_msg subscriber.service.internal_logger.log :info, "callback-delivery" do "message (ID #{rec_msg.message_id}, ackID #{rec_msg.ack_id}) delivery to user callbacks" end - @subscriber.callback.call rec_msg unless stopped? + @subscriber.callback.call rec_msg unless nack_immediately? rescue StandardError => e subscriber.service.internal_logger.log :info, "callback-exceptions" do "message (ID #{rec_msg.message_id}, ackID #{rec_msg.ack_id}) caused a user callback exception: " \ @@ -475,7 +514,7 @@ def perform_callback_sync rec_msg @subscriber.error! e ensure release rec_msg - if @sequencer && running? + if @sequencer && !nack_immediately? begin @sequencer.next rec_msg rescue OrderedMessageDeliveryError => e diff --git a/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/inventory_test.rb b/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/inventory_test.rb index 4d63d0530d5d..a77ccdbf9968 100644 --- a/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/inventory_test.rb +++ b/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/inventory_test.rb @@ -271,4 +271,66 @@ _(inventory.min_duration_per_lease_extension).must_equal 10 end + + it "waits until inventory is empty" do + subscriber_mock = Minitest::Mock.new + inventory = Google::Cloud::PubSub::MessageListener::Inventory.new subscriber_mock, + limit: 1000, + bytesize: 100_000, + extension: 3600, + max_duration_per_lease_extension: 0, + min_duration_per_lease_extension: 0 + + assert inventory.wait_until_empty(0.01) + + inventory.add rec_msg1_grpc + refute inventory.empty? + + thread = Thread.new do + inventory.remove "ack-id-1111" + end + + result = inventory.wait_until_empty 1.0 + thread.join + assert result + assert inventory.empty? + end + + it "returns false when wait_until_empty times out" do + subscriber_mock = Minitest::Mock.new + inventory = Google::Cloud::PubSub::MessageListener::Inventory.new subscriber_mock, + limit: 1000, + bytesize: 100_000, + extension: 3600, + max_duration_per_lease_extension: 0, + min_duration_per_lease_extension: 0 + + inventory.add rec_msg1_grpc + result = inventory.wait_until_empty 0.01 + refute result + refute inventory.empty? + end + + it "returns when stopped even if inventory is not empty" do + subscriber_mock = Minitest::Mock.new + inventory = Google::Cloud::PubSub::MessageListener::Inventory.new subscriber_mock, + limit: 1000, + bytesize: 100_000, + extension: 3600, + max_duration_per_lease_extension: 0, + min_duration_per_lease_extension: 0 + + inventory.add rec_msg1_grpc + refute inventory.empty? + + thread = Thread.new do + sleep 0.05 + inventory.stop + end + + result = inventory.wait_until_empty + thread.join + refute result + refute inventory.empty? + end end diff --git a/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/stream_test.rb b/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/stream_test.rb index 5f457b90405c..d4dd0349b4c0 100644 --- a/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/stream_test.rb +++ b/google-cloud-pubsub/test/google/cloud/pubsub/message_listener/stream_test.rb @@ -217,4 +217,131 @@ def stub.modify_ack_deadline subscription:, ack_ids:, ack_deadline_seconds: listener.stop listener.wait! end + + it "nacks unprocessed messages when stopped with nack_immediately" do + pull_res1 = Google::Cloud::PubSub::V1::StreamingPullResponse.new received_messages: [rec_msg1_grpc] + stub = StreamingPullStub.new [[pull_res1]] + subscriber.service.mocked_subscription_admin = stub + + message_received = Concurrent::Event.new + block_callback = Concurrent::Event.new + + listener = subscriber.listen streams: 1 do |_msg| + message_received.set + block_callback.wait + end + listener.instance_variable_set :@shutdown_behavior, :nack_immediately + + listener.start + message_received.wait + + listener.stream_pool.first.stop + block_callback.set + listener.buffer.stop + + # Verifies that exactly one 0-second ModifyAckDeadline (NACK) was dispatched. + assert_equal 1, stub.modify_ack_deadline_requests.count { |req| req[2] == 0 } + end + + it "waits for processing when stopped with wait_for_processing" do + pull_res1 = Google::Cloud::PubSub::V1::StreamingPullResponse.new received_messages: [rec_msg1_grpc] + stub = StreamingPullStub.new [[pull_res1]] + subscriber.service.mocked_subscription_admin = stub + + message_processed = Concurrent::Event.new + + listener = subscriber.listen streams: 1 do |_msg| + message_processed.set + end + + listener.start + message_processed.wait + + stream = listener.stream_pool.first + stream.stop + stream.wait! + listener.buffer.stop + + assert message_processed.set? + # Confirms that no 0-second ModifyAckDeadline (NACK) was dispatched. + assert_equal 0, stub.modify_ack_deadline_requests.count { |req| req[2] == 0 } + end + + it "processes subsequent ordered messages when stopped with wait_for_processing" do + ordered_msg1_grpc = Google::Cloud::PubSub::V1::ReceivedMessage.new( + rec_message_hash("msg-1", 1111).tap { |h| h[:message][:ordering_key] = "key1" } + ) + ordered_msg2_grpc = Google::Cloud::PubSub::V1::ReceivedMessage.new( + rec_message_hash("msg-2", 2222).tap { |h| h[:message][:ordering_key] = "key1" } + ) + pull_res = Google::Cloud::PubSub::V1::StreamingPullResponse.new( + received_messages: [ordered_msg1_grpc, ordered_msg2_grpc] + ) + stub = StreamingPullStub.new [[pull_res]] + subscriber.service.mocked_subscription_admin = stub + + first_message_received = Concurrent::Event.new + block_first_message = Concurrent::Event.new + processed_messages = [] + + ordered_subscriber = Google::Cloud::PubSub::Subscriber.from_grpc( + Google::Cloud::PubSub::V1::Subscription.new(sub_hash.merge(enable_message_ordering: true)), + pubsub.service + ) + + listener = ordered_subscriber.listen streams: 1 do |msg| + if msg.data == "msg-1" + first_message_received.set + block_first_message.wait + end + processed_messages << msg.data + msg.ack! + end + + listener.start + first_message_received.wait + + stream = listener.stream_pool.first + # Stop the stream while msg-1 is still running and msg-2 is queued in sequencer + stream.stop + block_first_message.set + + stream.wait! 2.0 + listener.buffer.stop + + assert_equal ["msg-1", "msg-2"], processed_messages + assert_equal 0, stub.modify_ack_deadline_requests.count { |req| req[2] == 0 } + end + + it "respects overall timeout budget across inventory and callback thread pool in wait!" do + pull_res1 = Google::Cloud::PubSub::V1::StreamingPullResponse.new received_messages: [rec_msg1_grpc] + stub = StreamingPullStub.new [[pull_res1]] + subscriber.service.mocked_subscription_admin = stub + + message_received = Concurrent::Event.new + block_callback = Concurrent::Event.new + + listener = subscriber.listen streams: 1 do |_msg| + message_received.set + block_callback.wait + end + + listener.start + message_received.wait + + stream = listener.stream_pool.first + stream.stop + + start_time = Process.clock_gettime Process::CLOCK_MONOTONIC + # Wait with a short timeout of 0.2s while callback is blocked + stream.wait! 0.2 + elapsed = Process.clock_gettime(Process::CLOCK_MONOTONIC) - start_time + + block_callback.set + listener.buffer.stop + + # Total wait time should be bounded near 0.2s, and well below 2x timeout (0.4s) + assert_operator elapsed, :>=, 0.15 + assert_operator elapsed, :<, 0.35 + end end