diff --git a/lib/async/pool/controller.rb b/lib/async/pool/controller.rb index fb81ef4..f01a063 100644 --- a/lib/async/pool/controller.rb +++ b/lib/async/pool/controller.rb @@ -166,16 +166,34 @@ def acquire # Make the resource resources and let waiting tasks know that there is something resources. def release(resource) - processed = false + unless usage = decrement_usage(resource) + return false + end - # A resource that is not good should also not be reusable. if resource.reusable? - processed = reuse(resource) + # If the resource was fully utilized, it now becomes available: + if usage == resource.concurrency - 1 + @available.push(resource) + end + else + # The resource must not be acquired again, but it cannot be retired until all existing users have released it: + @available.delete(resource) + + if usage.zero? + retire(resource) + resource = nil + return true + end end + @mutex.synchronize{@condition.broadcast} + # @policy.released(self, resource) + + resource = nil + return true ensure - retire(resource) unless processed + retire(resource) if resource end # Drain the pool, closing all resources. @@ -282,31 +300,6 @@ def availability_summary # @resources.count{|resource, usage| usage == 0} # end - def reuse(resource) - Console.debug(self){"Reuse #{resource}"} - - usage = @resources[resource] - - if usage.nil? - return false - end - - if usage.zero? - raise "Trying to reuse unacquired resource: #{resource}!" - end - - # If the resource was fully utilized, it now becomes available: - if usage == resource.concurrency - @available.push(resource) - end - - @resources[resource] = usage - 1 - - @mutex.synchronize{@condition.broadcast} - - return true - end - def wait_for_resource # If we fail to create a resource (below), we will end up waiting for one to become resources. until resource = available_resource @@ -345,12 +338,23 @@ def available_resource return resource rescue Exception - reuse(resource) if resource + release(resource) if resource raise end private + def decrement_usage(resource) + usage = @resources[resource] + return unless usage + + if usage.zero? + raise "Trying to reuse unacquired resource: #{resource}!" + end + + return @resources[resource] = usage - 1 + end + # Acquire an existing resource with zero usage. # If there are resources that are in use, wait until they are released. def acquire_existing_resource @@ -379,8 +383,8 @@ def acquire_or_create_resource return resource else - retire(resource) @available.pop + retire(resource) if usage.zero? end else # The resource has been removed already, so skip it and remove it from the availability list. diff --git a/test/async/pool/multiplex.rb b/test/async/pool/multiplex.rb index 93453d6..8f8f660 100644 --- a/test/async/pool/multiplex.rb +++ b/test/async/pool/multiplex.rb @@ -48,6 +48,112 @@ end end + with "a non-reusable resource" do + let(:constructor) {lambda{Async::Pool::Resource.new(3)}} + let(:pool) {subject.new(constructor, limit: 1)} + + it "retires the resource after the final release" do + resource1 = pool.acquire + resource2 = pool.acquire + + mock(resource1) do |mock| + mock.replace(:reusable?){false} + end + + pool.release(resource1) + + expect(pool.resources[resource1]).to be == 1 + expect(pool.available).to be == [] + expect(resource1).not.to be(:closed?) + + pool.release(resource2) + + expect(pool.resources).not.to be(:key?, resource1) + expect(resource1).to be(:closed?) + end + + it "continues to occupy the pool until the final release" do + resource1 = pool.acquire + resource2 = pool.acquire + + mock(resource1) do |mock| + mock.replace(:reusable?){false} + end + + pool.release(resource1) + + state = :waiting + acquired = nil + acquire_task = Async do + acquired = pool.acquire + state = :acquired + pool.release(acquired) + end + + expect(state).to be == :waiting + expect(pool.size).to be == 1 + + waited = false + wait_task = Async do + pool.wait_until_free + waited = true + end + + expect(waited).to be == false + + pool.release(resource2) + + acquire_task.wait + wait_task.wait + + expect(state).to be == :acquired + expect(acquired).not.to be_equal(resource1) + expect(waited).to be == true + end + + it "defers retirement when acquisition discovers an active non-viable resource" do + resource = pool.acquire + + mock(resource) do |mock| + mock.replace(:viable?){false} + mock.replace(:reusable?){false} + end + + state = :waiting + acquired = nil + acquire_task = Async do + acquired = pool.acquire + state = :acquired + pool.release(acquired) + end + + expect(state).to be == :waiting + expect(pool.resources[resource]).to be == 1 + expect(pool.available).to be == [] + expect(resource).not.to be(:closed?) + + pool.release(resource) + acquire_task.wait + + expect(state).to be == :acquired + expect(acquired).not.to be_equal(resource) + expect(resource).to be(:closed?) + end + + it "allows explicit retirement while the resource is active" do + resource1 = pool.acquire + resource2 = pool.acquire + + expect(pool.retire(resource1)).to be == true + + expect(pool.resources).not.to be(:key?, resource1) + expect(resource1).to be(:closed?) + + pool.release(resource1) + pool.release(resource2) + end + end + with "#prune" do it "removes the item from the availabilty list when it is retired" do object = pool.acquire