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
66 changes: 35 additions & 31 deletions lib/async/pool/controller.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
106 changes: 106 additions & 0 deletions test/async/pool/multiplex.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading