Packages
fixpoint
0.21.2
0.22.1
0.21.5
0.21.4
0.21.3
0.21.2
0.21.1
0.21.0
0.20.6
0.20.5
0.20.4
0.20.3
0.20.2
0.20.1
0.19.5
0.19.4
0.19.3
0.19.2
0.19.1
0.18.2
0.18.1
0.17.6
0.17.5
0.17.4
0.17.3
0.17.2
0.17.1
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.6
0.15.5
0.15.4
0.15.3
0.15.2
0.15.1
0.15.0
0.14.9
0.14.8
0.14.7
0.14.6
0.14.5
0.14.4
0.14.3
0.14.2
0.14.1
0.13.5
0.13.4
0.13.2
0.13.1
0.12.9
0.12.8
0.12.7
0.12.6
0.12.5
0.12.4
0.12.2
0.12.1
0.11.8
0.11.7
0.11.6
0.11.5
0.11.4
0.11.3
0.11.2
0.11.1
0.10.7
0.10.6
0.10.5
0.10.4
0.10.3
0.10.2
0.10.1
0.9.12
0.9.11
0.9.10
0.9.9
0.9.8
0.9.7
0.9.6
0.9.5
0.9.4
0.9.3
0.9.2
0.9.1
0.9.0
0.8.52
0.8.51
0.8.50
0.8.49
0.8.48
0.8.46
0.8.44
0.8.43
0.8.42
0.8.41
0.8.40
0.8.39
0.8.38
0.8.37
0.8.36
0.8.35
0.8.34
0.8.33
0.8.32
0.8.31
0.8.30
0.8.29
0.8.28
0.8.27
0.8.26
0.8.25
0.8.24
0.8.23
0.8.22
0.8.21
0.8.20
0.8.19
0.8.18
0.8.17
0.8.16
0.8.15
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.10
0.7.9
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.12
0.5.11
0.5.10
0.5.9
0.5.8
0.5.7
0.5.6
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
0.4.0
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.2.3
0.2.2
0.2.1
0.1.3
0.1.2
0.1.1
0.1.0
Constraint Programming Solver
Current section
Files
Jump to
Current section
Files
lib/solver/space/thread_pool.ex
defmodule CPSolver.Space.ThreadPool do
use GenServer
alias InPlace.Array
@impl true
def init(pool_size) do
pool_ref = Array.new(1, pool_size)
{:ok, %{space_queue: :queue.new(), pool_size: pool_size, pool_ref: pool_ref}}
end
@impl true
def handle_call(:checkout, {caller_pid, _ref} = caller, %{space_queue: queue} = state) do
if Process.alive?(caller_pid) do
if checkout_impl?(state) do
{:reply, true, state}
else
{:noreply, Map.put(state, :space_queue, :queue.in(caller, queue))}
end
else
{:noreply, state}
end
end
def handle_call(
:get_pool_state,
_caller,
%{pool_size: pool_size, pool_ref: pool_ref, space_queue: queue} = state
) do
{:reply, {:ok, %{queue: queue, pool_size: pool_size, available: get_free_threads(pool_ref)}},
state}
end
@impl true
def handle_cast(:checkin, state) do
updated_queue = checkin_impl(state)
{:noreply, Map.put(state, :space_queue, updated_queue)}
end
## API
def new(pool_size) when is_integer(pool_size) and pool_size > 0 do
{:ok, _pid} = GenServer.start(__MODULE__, pool_size)
end
def run_task(task, thread_pool, timeout \\ :infinity) when is_function(task) do
checkout(thread_pool, timeout)
try do
task.()
after
checkin(thread_pool)
end
end
def checkout(thread_pool, timeout \\ :infinity) when is_pid(thread_pool) do
GenServer.call(thread_pool, :checkout, timeout)
end
def checkin(thread_pool) when is_pid(thread_pool) do
GenServer.cast(thread_pool, :checkin)
end
def get_pool_state(thread_pool) do
GenServer.call(thread_pool, :get_pool_state)
end
defp get_free_threads(pool_ref) do
Array.get(pool_ref, 1)
end
defp checkout_impl?(%{pool_ref: pool_ref} = _state) do
case get_free_threads(pool_ref) do
free_threads when free_threads > 0 ->
Array.put(pool_ref, 1, free_threads - 1)
true
0 ->
## All threads are taken
false
_overspill ->
throw({:error, :thread_pool_checkout_error})
end
end
defp checkin_impl(%{pool_size: pool_size, pool_ref: pool_ref, space_queue: queue} = _state) do
increase_available_pool_count(pool_ref, pool_size)
{process_to_checkout, updated_queue} = get_waiting_process(queue)
if process_to_checkout do
decrease_available_pool_count(pool_ref)
## Wake up the waiting the process
GenServer.reply(process_to_checkout, true)
end
updated_queue
end
def get_waiting_process(queue) do
case :queue.out(queue) do
{:empty, q} ->
{nil, q}
{{:value, {pid, _ref} = caller}, q} ->
if Process.alive?(pid) do
{caller, q}
else
get_waiting_process(q)
end
end
end
defp increase_available_pool_count(pool_ref, pool_size) do
Array.update(pool_ref, 1, fn current ->
if current < pool_size do
current + 1
end
end)
end
defp decrease_available_pool_count(pool_ref) do
Array.update(pool_ref, 1, fn current ->
if current > 0 do
current - 1
end
end)
end
end