diff --git a/lib/async/limiter/generic.rb b/lib/async/limiter/generic.rb index 9184a47..4d45cff 100644 --- a/lib/async/limiter/generic.rb +++ b/lib/async/limiter/generic.rb @@ -6,6 +6,7 @@ require "async/task" require "async/deadline" +require "json" require_relative "timing/none" require_relative "timing/sliding_window" require_relative "token" @@ -135,6 +136,18 @@ def statistics end end + # Get a JSON-compatible representation of the limiter statistics. + # @returns [Hash] Statistics hash with current state. + def as_json(...) + statistics + end + + # Get a JSON string representation of the limiter statistics. + # @returns [String] JSON encoded statistics. + def to_json(...) + as_json.to_json(...) + end + protected def acquire_synchronized(timeout, cost, **options) diff --git a/lib/async/limiter/limited.rb b/lib/async/limiter/limited.rb index 059596b..6fbb79a 100644 --- a/lib/async/limiter/limited.rb +++ b/lib/async/limiter/limited.rb @@ -28,6 +28,8 @@ def initialize(limit = 1, timing: Timing::None, parent: nil) @limit = limit @count = 0 + @waiting_count = 0 + @reacquire_waiting_count = 0 @available = ConditionVariable.new end @@ -38,6 +40,26 @@ def initialize(limit = 1, timing: Timing::None, parent: nil) # @attribute [Integer] Current count of active tasks. attr_reader :count + # @returns [Integer] Current count of active tasks. + def acquired_count + @mutex.synchronize{@count} + end + + # @returns [Integer] Current count of available capacity. + def available_count + @mutex.synchronize{@limit - @count} + end + + # @returns [Integer] Current count of tasks waiting for capacity. + def waiting_count + @mutex.synchronize{@waiting_count} + end + + # @returns [Integer] Current count of reacquiring tasks waiting for capacity. + def reacquire_waiting_count + @mutex.synchronize{@reacquire_waiting_count} + end + # Check if a new task can be acquired. # @returns [Boolean] True if under the limit. def limited? @@ -64,6 +86,10 @@ def statistics { limit: @limit, count: @count, + acquired_count: @count, + available_count: @limit - @count, + waiting_count: @waiting_count, + reacquire_waiting_count: @reacquire_waiting_count, timing: @timing.statistics } end @@ -72,15 +98,23 @@ def statistics protected # Acquire resource with optional deadline. - def acquire_resource(deadline, **options) + def acquire_resource(deadline, reacquire: false, **options) # Fast path: immediate return for expired deadlines, but only if at capacity return nil if deadline&.expired? && @count >= @limit + waiting = false + # Wait for capacity with deadline tracking while @count >= @limit remaining = deadline&.remaining return nil if remaining && remaining <= 0 + unless waiting + @waiting_count += 1 + @reacquire_waiting_count += 1 if reacquire + waiting = true + end + unless @available.wait(@mutex, remaining) return nil # Timeout exceeded end @@ -89,6 +123,11 @@ def acquire_resource(deadline, **options) @count += 1 return true + ensure + if waiting + @waiting_count -= 1 + @reacquire_waiting_count -= 1 if reacquire + end end # Release resource. diff --git a/lib/async/limiter/queued.rb b/lib/async/limiter/queued.rb index a9da561..8486caa 100644 --- a/lib/async/limiter/queued.rb +++ b/lib/async/limiter/queued.rb @@ -33,11 +33,33 @@ def self.default_queue def initialize(queue = self.class.default_queue, timing: Timing::None, parent: nil) super(timing: timing, parent: parent) @queue = queue + @acquired_count = 0 + @reacquire_waiting_count = 0 end # @attribute [Queue] The queue managing resources. attr_reader :queue + # @returns [Integer] Current count of acquired resources. + def acquired_count + @mutex.synchronize{@acquired_count} + end + + # @returns [Integer] Current count of available resources. + def available_count + @queue.size + end + + # @returns [Integer] Current count of tasks waiting for resources. + def waiting_count + @queue.waiting_count + end + + # @returns [Integer] Current count of reacquiring tasks waiting for resources. + def reacquire_waiting_count + @mutex.synchronize{@reacquire_waiting_count} + end + # Check if a new task can be acquired. # @returns [Boolean] True if resources are available. def limited? @@ -60,6 +82,10 @@ def statistics { waiting: @queue.waiting_count, available: @queue.size, + acquired_count: @acquired_count, + available_count: @queue.size, + waiting_count: @queue.waiting_count, + reacquire_waiting_count: @reacquire_waiting_count, timing: @timing.statistics } end @@ -69,14 +95,23 @@ def statistics # Acquire a resource from the queue with optional deadline. def acquire_resource(deadline, reacquire: false, **options) + @reacquire_waiting_count += 1 if reacquire + @mutex.unlock - return @queue.pop(timeout: deadline&.remaining, **options) + resource = @queue.pop(timeout: deadline&.remaining, **options) + return resource ensure @mutex.lock + @reacquire_waiting_count -= 1 if reacquire + @acquired_count += 1 if resource end # Release a previously acquired resource back to the queue. def release_resource(value) + @mutex.synchronize do + @acquired_count -= 1 if @acquired_count > 0 + end + # Return a default resource to the queue: @queue.push(value) end diff --git a/test/async/limiter/limited.rb b/test/async/limiter/limited.rb index 27bd17b..219e22c 100644 --- a/test/async/limiter/limited.rb +++ b/test/async/limiter/limited.rb @@ -53,6 +53,43 @@ expect(limiter.acquire(timeout: 0)).to be == true end + it "tracks waiting tasks" do + limiter.acquire + limiter.acquire + + thread = Thread.new do + limiter.acquire + end + + Thread.pass until limiter.waiting_count == 1 + + expect(limiter.statistics[:waiting_count]).to be == 1 + expect(limiter.statistics[:reacquire_waiting_count]).to be == 0 + + limiter.release + expect(thread.value).to be == true + expect(limiter.statistics[:waiting_count]).to be == 0 + end + + it "tracks reacquiring waiting tasks separately" do + limiter.acquire + limiter.acquire + + thread = Thread.new do + limiter.acquire(reacquire: true) + end + + Thread.pass until limiter.reacquire_waiting_count == 1 + + expect(limiter.statistics[:waiting_count]).to be == 1 + expect(limiter.statistics[:reacquire_waiting_count]).to be == 1 + + limiter.release + expect(thread.value).to be == true + expect(limiter.statistics[:waiting_count]).to be == 0 + expect(limiter.statistics[:reacquire_waiting_count]).to be == 0 + end + it "handles deadline timeouts during condition variable waits" do # Fill limiter to capacity limiter.acquire # 1/2 diff --git a/test/async/limiter/queued.rb b/test/async/limiter/queued.rb index 9cefc28..3426bf8 100644 --- a/test/async/limiter/queued.rb +++ b/test/async/limiter/queued.rb @@ -102,6 +102,39 @@ # Should be empty again expect(limiter).to be(:limited?) end + + it "tracks acquired resources" do + limiter.expand(2, "tracked_resource") + + resource = limiter.acquire(timeout: 0) + + expect(resource).to be == "tracked_resource" + expect(limiter.statistics[:acquired_count]).to be == 1 + expect(limiter.statistics[:available_count]).to be == 1 + + limiter.release(resource) + + expect(limiter.statistics[:acquired_count]).to be == 0 + expect(limiter.statistics[:available_count]).to be == 2 + end + + it "tracks reacquiring waiting tasks separately" do + task = reactor.async do + limiter.acquire(reacquire: true) + end + + sleep 0.01 until limiter.reacquire_waiting_count == 1 + + expect(limiter.statistics[:waiting_count]).to be == 1 + expect(limiter.statistics[:reacquire_waiting_count]).to be == 1 + + limiter.release("resource") + + expect(task.wait).to be == "resource" + expect(limiter.statistics[:waiting_count]).to be == 0 + expect(limiter.statistics[:reacquire_waiting_count]).to be == 0 + expect(limiter.statistics[:acquired_count]).to be == 1 + end end with "priority queue" do diff --git a/test/async/limiter/statistics.rb b/test/async/limiter/statistics.rb index ac987f3..f764902 100644 --- a/test/async/limiter/statistics.rb +++ b/test/async/limiter/statistics.rb @@ -4,6 +4,7 @@ # Copyright, 2025, by Shopify Inc. require "async/limiter" +require "json" describe Async::Limiter::Generic do it "provides basic statistics" do @@ -25,6 +26,10 @@ expect(statistics).to be_a(Hash) expect(statistics[:limit]).to be == 3 expect(statistics[:count]).to be == 0 + expect(statistics[:acquired_count]).to be == 0 + expect(statistics[:available_count]).to be == 3 + expect(statistics[:waiting_count]).to be == 0 + expect(statistics[:reacquire_waiting_count]).to be == 0 expect(statistics[:timing]).to be_a(Hash) end @@ -36,6 +41,8 @@ expect(statistics[:limit]).to be == 3 expect(statistics[:count]).to be == 2 + expect(statistics[:acquired_count]).to be == 2 + expect(statistics[:available_count]).to be == 1 end it "updates statistics after release" do @@ -53,6 +60,14 @@ expect(limiter.statistics[:count]).to be == 0 end + it "provides a JSON representation" do + statistics = JSON.parse(limiter.to_json) + + expect(statistics["limit"]).to be == 3 + expect(statistics["acquired_count"]).to be == 0 + expect(statistics["available_count"]).to be == 3 + end + it "is thread-safe" do require "async" @@ -100,6 +115,10 @@ expect(statistics).to be_a(Hash) expect(statistics[:waiting]).to be_a(Integer) expect(statistics[:available]).to be_a(Integer) + expect(statistics[:acquired_count]).to be == 0 + expect(statistics[:available_count]).to be_a(Integer) + expect(statistics[:waiting_count]).to be_a(Integer) + expect(statistics[:reacquire_waiting_count]).to be == 0 expect(statistics[:timing]).to be_a(Hash) end @@ -111,21 +130,39 @@ # Initially empty expect(limiter.statistics[:available]).to be == 0 + expect(limiter.statistics[:available_count]).to be == 0 # Add resources limiter.release("worker1") limiter.release("worker2") expect(limiter.statistics[:available]).to be == 2 + expect(limiter.statistics[:available_count]).to be == 2 # Consume one resource resource = limiter.acquire(timeout: 0) expect(resource).to be == "worker1" expect(limiter.statistics[:available]).to be == 1 + expect(limiter.statistics[:available_count]).to be == 1 + expect(limiter.statistics[:acquired_count]).to be == 1 # Return resource limiter.release(resource) expect(limiter.statistics[:available]).to be == 2 + expect(limiter.statistics[:available_count]).to be == 2 + expect(limiter.statistics[:acquired_count]).to be == 0 + end + + it "provides a JSON representation" do + require "async/queue" + + queue = Async::Queue.new + limiter = Async::Limiter::Queued.new(queue) + statistics = JSON.parse(limiter.to_json) + + expect(statistics["acquired_count"]).to be == 0 + expect(statistics["available_count"]).to be == 0 + expect(statistics["waiting_count"]).to be == 0 end end