From 30b15f3e3fea94913239cb37d0e7a3ebe9fd75bc Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 31 Aug 2026 13:15:30 +1200 Subject: [PATCH 1/3] Use an ordered set for available resources Assisted-By: devx/618580b0-d55f-4c2e-95b6-87e648f60543 --- lib/async/pool/controller.rb | 22 ++++++++++------------ test/async/pool/multiplex.rb | 36 +++++++++++++++++++++++++++++------- 2 files changed, 39 insertions(+), 19 deletions(-) diff --git a/lib/async/pool/controller.rb b/lib/async/pool/controller.rb index f01a063..33cbc9a 100644 --- a/lib/async/pool/controller.rb +++ b/lib/async/pool/controller.rb @@ -12,6 +12,7 @@ require "async" require "async/semaphore" +require "set" require "thread" module Async @@ -44,9 +45,9 @@ def initialize(constructor, limit: nil, concurrency: 1, policy: nil, tags: nil) # All available resources: @resources = {} - # Resources which may be available to be acquired: - # This list may contain false positives, or resources which were okay but have since entered a state which is unusuable. - @available = [] + # Resources which may be available to be acquired. Resources are acquired in insertion order. Adding an existing resource preserves its position; a fully utilized resource is removed and reinserted at the end when capacity becomes available again. + # This set may contain false positives, or resources which were okay but have since entered a state which is unusable. + @available = Set.new # Used to signal when a resource has been released: @mutex = Thread::Mutex.new @@ -171,10 +172,7 @@ def release(resource) end if resource.reusable? - # If the resource was fully utilized, it now becomes available: - if usage == resource.concurrency - 1 - @available.push(resource) - end + @available.add(resource) else # The resource must not be acquired again, but it cannot be retired until all existing users have released it: @available.delete(resource) @@ -319,7 +317,7 @@ def create_resource # Make the resource available if it can be used multiple times: if resource.concurrency > 1 - @available.push(resource) + @available.add(resource) end end @@ -371,24 +369,24 @@ def acquire_existing_resource end def acquire_or_create_resource - while resource = @available.last + while resource = @available.first if usage = @resources[resource] and usage < resource.concurrency if resource.viable? usage = (@resources[resource] += 1) if usage == resource.concurrency # The resource is used up to it's limit: - @available.pop + @available.delete(resource) end return resource else - @available.pop + @available.delete(resource) retire(resource) if usage.zero? end else # The resource has been removed already, so skip it and remove it from the availability list. - @available.pop + @available.delete(resource) end end diff --git a/test/async/pool/multiplex.rb b/test/async/pool/multiplex.rb index 8f8f660..a8644d4 100644 --- a/test/async/pool/multiplex.rb +++ b/test/async/pool/multiplex.rb @@ -26,7 +26,7 @@ pool.release(object) expect(pool).to be(:active?) - expect(pool.available).to be == [object] + expect(pool.available.to_a).to be == [object] end it "can acquire and release the same object up to the concurrency limit" do @@ -41,10 +41,32 @@ expect(pool.available).to be(:empty?) pool.release(object1) - expect(pool.available).to be == [object1] + expect(pool.available.to_a).to be == [object1] pool.release(object2) - expect(pool.available).to be == [object1] + expect(pool.available.to_a).to be == [object1] + end + + it "acquires resources in insertion order" do + resource1 = pool.acquire + resource2 = pool.acquire + resource3 = pool.acquire + + expect(resource2).to be_equal(resource1) + expect(resource3).not.to be_equal(resource1) + + pool.release(resource1) + + expect(pool.available.to_a).to be == [resource3, resource1] + + pool.release(resource2) + expect(pool.available.to_a).to be == [resource3, resource1] + + pool.acquire do |resource| + expect(resource).to be_equal(resource3) + end + + pool.release(resource3) end end @@ -63,7 +85,7 @@ pool.release(resource1) expect(pool.resources[resource1]).to be == 1 - expect(pool.available).to be == [] + expect(pool.available).to be(:empty?) expect(resource1).not.to be(:closed?) pool.release(resource2) @@ -129,7 +151,7 @@ expect(state).to be == :waiting expect(pool.resources[resource]).to be == 1 - expect(pool.available).to be == [] + expect(pool.available).to be(:empty?) expect(resource).not.to be(:closed?) pool.release(resource) @@ -166,7 +188,7 @@ pool.prune - expect(pool.available).to be == [] + expect(pool.available).to be(:empty?) end it "puts the item back into the available list if it is reusable" do @@ -177,7 +199,7 @@ pool.prune - expect(pool.available).to be == [object] + expect(pool.available.to_a).to be == [object] end end end From e27a8a53347fa2ba455e7cd401c9d1eb26404534 Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 31 Aug 2026 13:18:11 +1200 Subject: [PATCH 2/3] Compare pooled resources by identity Assisted-By: devx/618580b0-d55f-4c2e-95b6-87e648f60543 --- lib/async/pool/controller.rb | 8 ++++---- test/async/pool/controller.rb | 31 +++++++++++++++++++++++++++++++ 2 files changed, 35 insertions(+), 4 deletions(-) diff --git a/lib/async/pool/controller.rb b/lib/async/pool/controller.rb index 33cbc9a..3fb1b24 100644 --- a/lib/async/pool/controller.rb +++ b/lib/async/pool/controller.rb @@ -42,12 +42,12 @@ def initialize(constructor, limit: nil, concurrency: 1, policy: nil, tags: nil) @tags = tags - # All available resources: - @resources = {} + # All allocated resources. Each resource is tracked by identity, irrespective of any value equality it may define: + @resources = {}.compare_by_identity - # Resources which may be available to be acquired. Resources are acquired in insertion order. Adding an existing resource preserves its position; a fully utilized resource is removed and reinserted at the end when capacity becomes available again. + # Resources which may be available to be acquired. Resources are compared by identity and acquired in insertion order. Adding an existing resource preserves its position; a fully utilized resource is removed and reinserted at the end when capacity becomes available again. # This set may contain false positives, or resources which were okay but have since entered a state which is unusable. - @available = Set.new + @available = Set.new.compare_by_identity # Used to signal when a resource has been released: @mutex = Thread::Mutex.new diff --git a/test/async/pool/controller.rb b/test/async/pool/controller.rb index 25708d7..a4b6889 100644 --- a/test/async/pool/controller.rb +++ b/test/async/pool/controller.rb @@ -53,6 +53,37 @@ end end + with "resources which compare equal" do + let(:resource_class) do + Class.new(Async::Pool::Resource) do + def == other + other.instance_of?(self.class) + end + + alias eql? == + + def hash + self.class.hash + end + end + end + + let(:pool) {subject.new(resource_class)} + + it "tracks each resource by identity" do + resource1 = pool.acquire + resource2 = pool.acquire + + expect(resource2).not.to be_equal(resource1) + expect(pool.resources.size).to be == 2 + + pool.release(resource1) + pool.release(resource2) + + expect(pool.available.to_a).to be == [resource1, resource2] + end + end + with "a limited pool" do let(:pool) {subject.new(Async::Pool::Resource, limit: 1)} From 4539ee837207b80d7985827bde85d58887b871bc Mon Sep 17 00:00:00 2001 From: Samuel Williams Date: Mon, 31 Aug 2026 13:34:47 +1200 Subject: [PATCH 3/3] Add unreleased note for ordered resources Assisted-By: devx/618580b0-d55f-4c2e-95b6-87e648f60543 --- releases.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/releases.md b/releases.md index 8fa3b57..dc7bf16 100644 --- a/releases.md +++ b/releases.md @@ -1,3 +1,7 @@ # Releases +## Unreleased + +- Use identity-based tracking for pooled resources and preserve insertion order when selecting available resources. + ## v0.11.2