Packages
phoenix
1.7.23
1.8.9
1.8.8
1.8.7
1.8.6
1.8.5
1.8.4
1.8.3
1.8.2
1.8.1
1.8.0
1.8.0-rc.4
1.8.0-rc.3
1.8.0-rc.2
1.8.0-rc.1
1.8.0-rc.0
1.7.24
1.7.23
1.7.22
1.7.21
1.7.20
1.7.19
1.7.18
1.7.17
1.7.16
1.7.15
1.7.14
1.7.13
1.7.12
1.7.11
1.7.10
1.7.9
1.7.8
1.7.7
1.7.6
1.7.5
1.7.4
1.7.3
1.7.2
1.7.1
1.7.0
1.7.0-rc.3
1.7.0-rc.2
1.7.0-rc.1
1.7.0-rc.0
1.6.17
1.6.16
1.6.15
1.6.14
1.6.13
1.6.12
1.6.11
1.6.10
1.6.9
1.6.8
1.6.7
1.6.6
1.6.5
1.6.4
1.6.3
1.6.2
1.6.1
1.6.0
1.6.0-rc.1
1.6.0-rc.0
1.5.15
1.5.14
1.5.13
1.5.12
1.5.11
1.5.10
1.5.9
1.5.8
1.5.7
1.5.6
1.5.5
1.5.4
1.5.3
1.5.2
1.5.1
1.5.0
1.5.0-rc.0
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.12
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.4.0-rc.3
1.4.0-rc.2
1.4.0-rc.1
1.4.0-rc.0
1.3.5
1.3.4
1.3.3
1.3.2
1.3.1
1.3.0
1.3.0-rc.3
1.3.0-rc.2
1.3.0-rc.1
1.3.0-rc.0
1.2.5
1.2.4
1.2.3
1.2.2
1.2.1
1.2.0
1.2.0-rc.1
1.2.0-rc.0
1.1.9
1.1.8
1.1.7
1.1.6
1.1.5
1.1.4
1.1.3
1.1.2
1.1.1
1.1.0
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.17.1
0.17.0
0.16.1
0.16.0
0.15.0
0.14.0
0.13.1
0.13.0
0.12.0
0.11.0
0.10.0
0.9.0
0.8.0
0.7.2
0.7.1
0.7.0
0.6.2
0.6.1
0.6.0
0.5.0
0.4.1
0.4.0
0.3.1
0.3.0
0.2.11
0.2.10
0.2.9
0.2.8
0.2.7
0.2.6
0.2.5
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.0
Productive. Reliable. Fast. A productive web framework that does not compromise speed or maintainability.
Security advisory:
This version has known vulnerabilities.
View advisories
Current section
Files
Jump to
Current section
Files
lib/phoenix/socket/pool_supervisor.ex
defmodule Phoenix.Socket.PoolSupervisor do
@moduledoc false
use Supervisor
def start_link({endpoint, name, partitions}) do
Supervisor.start_link(
__MODULE__,
{endpoint, name, partitions},
name: Module.concat(endpoint, name)
)
end
def start_child(socket, key, spec) do
%{endpoint: endpoint, handler: name} = socket
case endpoint.config({:socket, name}) do
ets when not is_nil(ets) ->
partitions = :ets.lookup_element(ets, :partitions, 2)
sup = :ets.lookup_element(ets, :erlang.phash2(key, partitions), 2)
DynamicSupervisor.start_child(sup, spec)
nil ->
raise ArgumentError, """
no socket supervision tree found for #{inspect(name)}.
Ensure your #{inspect(endpoint)} contains a socket mount, for example:
socket "/socket", #{inspect(name)},
websocket: true,
longpoll: true
"""
end
end
def start_pooled(ref, i) do
case DynamicSupervisor.start_link(strategy: :one_for_one) do
{:ok, pid} ->
:ets.insert(ref, {i, pid})
{:ok, pid}
{:error, reason} ->
{:error, reason}
end
end
@impl true
def init({endpoint, name, partitions}) do
# TODO: Use persisent term on Elixir v1.12+
ref = :ets.new(name, [:public, read_concurrency: true])
:ets.insert(ref, {:partitions, partitions})
Phoenix.Config.permanent(endpoint, {:socket, name}, ref)
children =
for i <- 0..(partitions - 1) do
%{
id: i,
start: {__MODULE__, :start_pooled, [ref, i]},
type: :supervisor,
shutdown: :infinity
}
end
Supervisor.init(children, strategy: :one_for_one)
end
end
defmodule Phoenix.Socket.PoolDrainer do
@moduledoc false
use GenServer
require Logger
def child_spec({_endpoint, name, opts} = tuple) do
# The process should terminate within shutdown but,
# in case it doesn't, we will be killed if we exceed
# double of that
%{
id: {:terminator, name},
start: {__MODULE__, :start_link, [tuple]},
shutdown: Keyword.get(opts[:drainer], :shutdown, 30_000)
}
end
def start_link(tuple) do
GenServer.start_link(__MODULE__, tuple)
end
@impl true
def init({endpoint, name, opts}) do
Process.flag(:trap_exit, true)
size = Keyword.get(opts[:drainer], :batch_size, 10_000)
interval = Keyword.get(opts[:drainer], :batch_interval, 2_000)
log_level = Keyword.get(opts[:drainer], :log, opts[:log] || :info)
{:ok, {endpoint, name, size, interval, log_level}}
end
@impl true
def terminate(_reason, {endpoint, name, size, interval, log_level}) do
ets = endpoint.config({:socket, name})
partitions = :ets.lookup_element(ets, :partitions, 2)
{collection, total} =
Enum.map_reduce(0..(partitions - 1), 0, fn index, total ->
try do
sup = :ets.lookup_element(ets, index, 2)
children = DynamicSupervisor.which_children(sup)
{Enum.map(children, &elem(&1, 1)), total + length(children)}
catch
_, _ -> {[], total}
end
end)
rounds = div(total, size) + 1
for {pids, index} <-
collection |> Stream.concat() |> Stream.chunk_every(size) |> Stream.with_index(1) do
count = if index == rounds, do: length(pids), else: size
:telemetry.execute(
[:phoenix, :socket_drain],
%{count: count, total: total, index: index, rounds: rounds},
%{
endpoint: endpoint,
socket: name,
interval: interval,
log: log_level
}
)
spawn(fn ->
for pid <- pids do
send(pid, %Phoenix.Socket.Broadcast{event: "phx_drain"})
end
end)
if index < rounds do
Process.sleep(interval)
end
end
end
end