Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions lib/async/limiter/generic.rb
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

require "async/task"
require "async/deadline"
require "json"
require_relative "timing/none"
require_relative "timing/sliding_window"
require_relative "token"
Expand Down Expand Up @@ -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)
Expand Down
41 changes: 40 additions & 1 deletion lib/async/limiter/limited.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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?
Expand All @@ -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
Expand All @@ -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
Expand All @@ -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.
Expand Down
37 changes: 36 additions & 1 deletion lib/async/limiter/queued.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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?
Expand All @@ -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
Expand All @@ -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
Expand Down
37 changes: 37 additions & 0 deletions test/async/limiter/limited.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
33 changes: 33 additions & 0 deletions test/async/limiter/queued.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
37 changes: 37 additions & 0 deletions test/async/limiter/statistics.rb
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
# Copyright, 2025, by Shopify Inc.

require "async/limiter"
require "json"

describe Async::Limiter::Generic do
it "provides basic statistics" do
Expand All @@ -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

Expand All @@ -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
Expand All @@ -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"

Expand Down Expand Up @@ -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

Expand All @@ -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

Expand Down
Loading